diff --git a/apps/backend/src/repositories/kuaishou-cloud-task-state-repo.ts b/apps/backend/src/repositories/kuaishou-cloud-task-state-repo.ts index 539e3d95..a605a945 100644 --- a/apps/backend/src/repositories/kuaishou-cloud-task-state-repo.ts +++ b/apps/backend/src/repositories/kuaishou-cloud-task-state-repo.ts @@ -1,4 +1,5 @@ import { query } from '../db/client.js' +import { TASK_STATUS } from '../domain/task-status.js' import { maskCode } from '../utils/masking.js' import { parseTaskContext } from '../utils/task-json.js' import { normalizeTimestampIso } from '../utils/time.js' @@ -6,6 +7,12 @@ import { normalizeTimestampIso } from '../utils/time.js' type QueryExecutor = (text: string, params?: unknown[]) => Promise<{ rows: unknown[] }> type JsonRecord = Record +export type KuaishouCloudSourceLoadStat = { + sourceKey: string + activeCount: number + redeemingCount: number +} + export type KuaishouCloudTaskStateSyncInput = { taskId: number | string executorKey?: unknown @@ -86,6 +93,49 @@ export async function syncKuaishouCloudTaskStateForTask( ) } +export async function listKuaishouCloudSourceLoadStats( + sourceKeys: unknown[] = [], + executor: QueryExecutor = query, +): Promise { + const candidates = normalizeStringArray(sourceKeys) + if (candidates.length === 0) { + return [] + } + + const activeStatuses = [ + TASK_STATUS.WAITING_BINDING, + TASK_STATUS.ROLE_CONFIRMED, + TASK_STATUS.REDEEMING, + TASK_STATUS.DISPATCHED_PENDING_RETURN, + ] + const result = await executor( + ` + SELECT + kcts.source_key, + COUNT(*) FILTER (WHERE ft.task_status = ANY($2::text[])) AS active_count, + COUNT(*) FILTER (WHERE ft.task_status = $3) AS redeeming_count + FROM fulfillment_tasks ft + JOIN kuaishou_cloud_task_states kcts ON kcts.task_id = ft.id + WHERE ft.executor_key = 'kuaishou_ct_assisted' + AND kcts.source_key = ANY($1::text[]) + AND ft.task_status = ANY($2::text[]) + GROUP BY kcts.source_key + `, + [candidates, activeStatuses, TASK_STATUS.REDEEMING], + ) + + return result.rows + .map((row) => { + const source = asRecord(row) + return { + sourceKey: String(source.source_key || '').trim(), + activeCount: Number(source.active_count || 0) || 0, + redeemingCount: Number(source.redeeming_count || 0) || 0, + } + }) + .filter((row) => row.sourceKey) +} + function resolveKuaishouCloudFlow(contextJson: unknown): JsonRecord | null { const context = parseTaskContext({ context_json: contextJson }) const flow = asRecord(context.kuaishouCloudFulfillment) @@ -93,13 +143,32 @@ function resolveKuaishouCloudFlow(contextJson: unknown): JsonRecord | null { } function asRecord(value: unknown): JsonRecord { - return value && typeof value === 'object' && !Array.isArray(value) ? value as JsonRecord : {} + return value && typeof value === 'object' && !Array.isArray(value) ? (value as JsonRecord) : {} } function normalizeStatus(value: unknown): string { return String(value || 'pending').trim() || 'pending' } +function normalizeStringArray(value: unknown): string[] { + if (Array.isArray(value)) { + return Array.from(new Set(value.map((item) => String(item || '').trim()).filter(Boolean))) + } + + if (typeof value === 'string') { + return Array.from( + new Set( + value + .split(',') + .map((item) => item.trim()) + .filter(Boolean), + ), + ) + } + + return [] +} + function pickFirstNonEmpty(values: unknown[]): string { for (const value of values) { const normalized = String(value || '').trim() diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/account-selector.test.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/account-selector.test.ts new file mode 100644 index 00000000..7ead69ee --- /dev/null +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/account-selector.test.ts @@ -0,0 +1,124 @@ +import assert from 'node:assert/strict' +import test from 'node:test' + +import { selectCloudtentaclesSourceForFulfillment } from './account-selector.js' + +const contexts = { + 'account-a': { + baseUrl: 'https://cloud.example.com', + token: 'token-a', + deviceId: '-', + deviceType: 0, + resolvedSourceKey: 'account-a', + }, + 'account-b': { + baseUrl: 'https://cloud.example.com', + token: 'token-b', + deviceId: '-', + deviceType: 0, + resolvedSourceKey: 'account-b', + }, +} + +test('selectCloudtentaclesSourceForFulfillment 已固定账号时不按压力切换', async () => { + const selection = await selectCloudtentaclesSourceForFulfillment( + { + binding: { + resolvedSourceKey: 'account-a', + cloudSourceKeys: ['account-a', 'account-b'], + }, + }, + { + resolveContextBySourceKeys: (sourceKeys) => { + const key = String(sourceKeys[0] || '').trim() + return contexts[key as keyof typeof contexts] + }, + listSourceLoadStats: async () => [ + { sourceKey: 'account-a', activeCount: 99, redeemingCount: 99 }, + { sourceKey: 'account-b', activeCount: 0, redeemingCount: 0 }, + ], + }, + ) + + assert.equal(selection.sourceKey, 'account-a') + assert.equal(selection.fixed, true) +}) + +test('selectCloudtentaclesSourceForFulfillment 多候选时选择压力最低账号', async () => { + const selection = await selectCloudtentaclesSourceForFulfillment( + { + binding: { + cloudSourceKeys: ['account-a', 'account-b'], + }, + }, + { + resolveContextBySourceKeys: (sourceKeys) => { + const key = String(sourceKeys[0] || '').trim() + return contexts[key as keyof typeof contexts] + }, + listSourceLoadStats: async () => [ + { sourceKey: 'account-a', activeCount: 3, redeemingCount: 1 }, + { sourceKey: 'account-b', activeCount: 0, redeemingCount: 0 }, + ], + }, + ) + + try { + assert.equal(selection.sourceKey, 'account-b') + assert.equal(selection.fixed, false) + } finally { + selection.release() + } +}) + +test('selectCloudtentaclesSourceForFulfillment 会跳过不可用账号', async () => { + const selection = await selectCloudtentaclesSourceForFulfillment( + { + binding: { + cloudSourceKeys: ['account-a', 'account-b'], + }, + }, + { + resolveContextBySourceKeys: (sourceKeys) => { + const key = String(sourceKeys[0] || '').trim() + if (key === 'account-a') { + throw new Error('账号已停用') + } + return contexts[key as keyof typeof contexts] + }, + listSourceLoadStats: async () => [], + }, + ) + + try { + assert.equal(selection.sourceKey, 'account-b') + assert.equal(selection.candidateCount, 1) + } finally { + selection.release() + } +}) + +test('selectCloudtentaclesSourceForFulfillment 用进程内占位分散并发准备', async () => { + const deps = { + resolveContextBySourceKeys: (sourceKeys: unknown[] = []) => { + const key = String(sourceKeys[0] || '').trim() + return contexts[key as keyof typeof contexts] + }, + listSourceLoadStats: async () => [], + } + const flow = { + binding: { + cloudSourceKeys: ['account-a', 'account-b'], + }, + } + const first = await selectCloudtentaclesSourceForFulfillment(flow, deps) + const second = await selectCloudtentaclesSourceForFulfillment(flow, deps) + + try { + assert.equal(first.sourceKey, 'account-a') + assert.equal(second.sourceKey, 'account-b') + } finally { + first.release() + second.release() + } +}) diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/account-selector.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/account-selector.ts new file mode 100644 index 00000000..5165f2ed --- /dev/null +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/account-selector.ts @@ -0,0 +1,182 @@ +import { listKuaishouCloudSourceLoadStats } from '../../../repositories/kuaishou-cloud-task-state-repo.js' +import { normalizeKuaishouCloudFlow, normalizeStringArray, type JsonObject } from './domain.js' +import { resolvePersistedCloudtentaclesContextBySourceKeys } from './cloudtentacles-context.js' + +type CloudtentaclesContext = ReturnType + +type CloudtentaclesSourceSelectionDeps = { + resolveContextBySourceKeys?: typeof resolvePersistedCloudtentaclesContextBySourceKeys + listSourceLoadStats?: typeof listKuaishouCloudSourceLoadStats +} + +export type CloudtentaclesSourceSelection = { + context: CloudtentaclesContext + sourceKey: string + candidateCount: number + activeCount: number + redeemingCount: number + inProcessCount: number + score: number + fixed: boolean + release: () => void +} + +const inProcessSelections = new Map() +let selectionLockTail: Promise = Promise.resolve() + +export async function selectCloudtentaclesSourceForFulfillment( + flowLike: unknown, + deps: CloudtentaclesSourceSelectionDeps = {}, +): Promise { + const flow = normalizeKuaishouCloudFlow(flowLike) + const resolveContext = + deps.resolveContextBySourceKeys || resolvePersistedCloudtentaclesContextBySourceKeys + const listSourceLoadStats = deps.listSourceLoadStats || listKuaishouCloudSourceLoadStats + const fixedSourceKey = String(flow.binding.resolvedSourceKey || '').trim() + const candidates = normalizeCloudtentaclesSourceCandidates( + fixedSourceKey + ? [fixedSourceKey, ...flow.binding.cloudSourceKeys] + : flow.binding.cloudSourceKeys, + ) + + if (fixedSourceKey) { + const context = resolveContext([fixedSourceKey, ...flow.binding.cloudSourceKeys]) + return { + context, + sourceKey: context.resolvedSourceKey, + candidateCount: candidates.length, + activeCount: 0, + redeemingCount: 0, + inProcessCount: 0, + score: 0, + fixed: true, + release: () => {}, + } + } + + return withCloudtentaclesSourceSelectionLock(async () => { + const contexts = [] + let lastError: unknown = null + + for (const sourceKey of candidates) { + try { + contexts.push(resolveContext([sourceKey])) + } catch (error) { + lastError = error + } + } + + if (contexts.length === 0) { + if (lastError) { + throw lastError + } + resolveContext(candidates) + throw new Error('所有 cloudtentacles 账号均不可用') + } + + const stats = await listSourceLoadStats(contexts.map((context) => context.resolvedSourceKey)) + const statsBySourceKey = new Map(stats.map((stat) => [stat.sourceKey, stat])) + const ranked = contexts + .map((context, index) => { + const stat = statsBySourceKey.get(context.resolvedSourceKey) || { + sourceKey: context.resolvedSourceKey, + activeCount: 0, + redeemingCount: 0, + } + const activeCount = normalizeCount(stat.activeCount) + const redeemingCount = normalizeCount(stat.redeemingCount) + const inProcessCount = normalizeCount(inProcessSelections.get(context.resolvedSourceKey)) + const score = (activeCount + inProcessCount) * 10 + redeemingCount * 20 + + return { + context, + sourceKey: context.resolvedSourceKey, + index, + activeCount, + redeemingCount, + inProcessCount, + score, + } + }) + .sort( + (left, right) => + left.score - right.score || + left.activeCount - right.activeCount || + left.redeemingCount - right.redeemingCount || + left.index - right.index || + left.sourceKey.localeCompare(right.sourceKey), + ) + + const selected = ranked[0] + if (!selected) { + throw new Error('所有 cloudtentacles 账号均不可用') + } + + reserveCloudtentaclesSource(selected.sourceKey) + + return { + context: selected.context, + sourceKey: selected.sourceKey, + candidateCount: contexts.length, + activeCount: selected.activeCount, + redeemingCount: selected.redeemingCount, + inProcessCount: selected.inProcessCount, + score: selected.score, + fixed: false, + release: () => releaseCloudtentaclesSource(selected.sourceKey), + } + }) +} + +export function buildCloudtentaclesSourceSelectionEventDetail( + selection: CloudtentaclesSourceSelection, +): JsonObject { + return { + sourceKey: selection.sourceKey, + candidateCount: selection.candidateCount, + activeCount: selection.activeCount, + redeemingCount: selection.redeemingCount, + inProcessCount: selection.inProcessCount, + score: selection.score, + fixed: selection.fixed, + } +} + +function normalizeCloudtentaclesSourceCandidates(sourceKeys: unknown[] = []) { + return Array.from(new Set(normalizeStringArray(sourceKeys))) +} + +async function withCloudtentaclesSourceSelectionLock(fn: () => Promise): Promise { + let releaseLock: () => void = () => {} + const waitForTurn = selectionLockTail + selectionLockTail = new Promise((resolve) => { + releaseLock = resolve + }) + + await waitForTurn + + try { + return await fn() + } finally { + releaseLock() + } +} + +function reserveCloudtentaclesSource(sourceKey: string) { + inProcessSelections.set(sourceKey, normalizeCount(inProcessSelections.get(sourceKey)) + 1) +} + +function releaseCloudtentaclesSource(sourceKey: string) { + const nextCount = normalizeCount(inProcessSelections.get(sourceKey)) - 1 + if (nextCount > 0) { + inProcessSelections.set(sourceKey, nextCount) + return + } + + inProcessSelections.delete(sourceKey) +} + +function normalizeCount(value: unknown) { + const count = Number(value || 0) || 0 + return Number.isFinite(count) && count > 0 ? count : 0 +} diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/index.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/index.ts index c9b52091..af2375f6 100644 --- a/apps/backend/src/services/fulfillment/kuaishou-cloud/index.ts +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/index.ts @@ -1,22 +1,18 @@ -import { createTaskEvent } from "../../../repositories/task-event-repo.js"; -import { updateTask } from "../../../repositories/task-repo.js"; -import { buildClaimUrl, createTaskClaimToken } from "../../claim/claim-service.js"; -import { - listCloudtentaclesSku, -} from "../../platforms/cloudtentacles/catalog-service.js"; -import { getCloudtentaclesKnapsack } from "../../platforms/cloudtentacles/knapsack-service.js"; +import { createTaskEvent } from '../../../repositories/task-event-repo.js' +import { updateTask } from '../../../repositories/task-repo.js' +import { buildClaimUrl, createTaskClaimToken } from '../../claim/claim-service.js' +import { listCloudtentaclesSku } from '../../platforms/cloudtentacles/catalog-service.js' +import { getCloudtentaclesKnapsack } from '../../platforms/cloudtentacles/knapsack-service.js' import { backCloudtentaclesVirtualNumber, getCloudtentaclesBindInfo, probeCloudtentaclesBindUrl, -} from "../../platforms/cloudtentacles/virtual-number-service.js"; -import { resolveCloudtentaclesConfig } from "../../platforms/cloudtentacles/helpers.js"; -import { - notifyKuaishouCloudBindUrlRefreshFailed, -} from "../../notification/domain-notifications.js"; -import { createHttpError } from "../../../utils/http.js"; -import { nowIso } from "../../../utils/time.js"; -import { TASK_STATUS, normalizeTaskStatus } from "../../../domain/task-status.js"; +} from '../../platforms/cloudtentacles/virtual-number-service.js' +import { resolveCloudtentaclesConfig } from '../../platforms/cloudtentacles/helpers.js' +import { notifyKuaishouCloudBindUrlRefreshFailed } from '../../notification/domain-notifications.js' +import { createHttpError } from '../../../utils/http.js' +import { nowIso } from '../../../utils/time.js' +import { TASK_STATUS, normalizeTaskStatus } from '../../../domain/task-status.js' import { KUAISHOU_CLOUD_FIXED_VN_KEY, isKuaishouCloudBindUrlFresh, @@ -26,32 +22,34 @@ import { normalizeKuaishouCloudRoleInfo, resolveKuaishouCloudBindUrlExpiresAt, type JsonObject, -} from "./domain.js"; +} from './domain.js' import { prepareKuaishouCloudBindResourceWithFallback, resolveKuaishouCloudBindingResources, resolveKuaishouCloudVnKeyCandidates, -} from "./binding-resources.js"; +} from './binding-resources.js' +import { resolvePersistedCloudtentaclesContextBySourceKeys } from './cloudtentacles-context.js' import { - resolvePersistedCloudtentaclesContextBySourceKeys, -} from "./cloudtentacles-context.js"; + buildCloudtentaclesSourceSelectionEventDetail, + selectCloudtentaclesSourceForFulfillment, +} from './account-selector.js' import { getTaskClaimExpiresAt, isClaimExpired, normalizeActor, parseTaskContext, -} from "./task-context.js"; -import type { TaskRow } from "../../../types/repository/rows.js"; +} from './task-context.js' +import type { TaskRow } from '../../../types/repository/rows.js' type CloudDeliveryPlanItem = { - cloudSkuId: number; - cloudSkuName: string; - quantity: number; - skuItem: JsonObject | null; - knapsackItem: JsonObject | null; - knapsackCount: number; - missingCount: number; -}; + cloudSkuId: number + cloudSkuName: string + quantity: number + skuItem: JsonObject | null + knapsackItem: JsonObject | null + knapsackCount: number + missingCount: number +} export { isKuaishouCloudBindUrlFresh, @@ -61,77 +59,70 @@ export { normalizeKuaishouCloudFlow, normalizeKuaishouCloudRoleInfo, resolveKuaishouCloudBindUrlExpiresAt, -} from "./domain.js"; -export { - resolvePersistedCloudtentaclesContextBySourceKeys, -} from "./cloudtentacles-context.js"; +} from './domain.js' +export { resolvePersistedCloudtentaclesContextBySourceKeys } from './cloudtentacles-context.js' export async function ensureTaskClaimLink(task: TaskRow) { - const tokenStatus = String(task?.primary_claim_token_status || "").trim(); - const token = String( - task?.primary_claim_token || task?.claim_token || "" - ).trim(); - const expiredAt = getTaskClaimExpiresAt(task); + const tokenStatus = String(task?.primary_claim_token_status || '').trim() + const token = String(task?.primary_claim_token || task?.claim_token || '').trim() + const expiredAt = getTaskClaimExpiresAt(task) - if (tokenStatus === "active" && token && !isClaimExpired(expiredAt)) { + if (tokenStatus === 'active' && token && !isClaimExpired(expiredAt)) { return { token, expiredAt, claimUrl: buildClaimUrl(token), - }; + } } - const claimToken = await createTaskClaimToken(task.id); + const claimToken = await createTaskClaimToken(task.id) return { token: claimToken.token, expiredAt: claimToken.expired_at, claimUrl: claimToken.claimUrl, - }; + } } -export async function prepareKuaishouCloudFulfillmentTask( - task: TaskRow, - options: JsonObject = {} -) { +export async function prepareKuaishouCloudFulfillmentTask(task: TaskRow, options: JsonObject = {}) { if (!isKuaishouCloudTask(task)) { - throw createHttpError("当前任务不是快手 Cloud 履约任务", { + throw createHttpError('当前任务不是快手 Cloud 履约任务', { statusCode: 409, - errorCode: "kuaishou_cloud_task_invalid", - }); + errorCode: 'kuaishou_cloud_task_invalid', + }) } - const now = nowIso(); - const actor = normalizeActor(options.actor); - const force = options.force === true; - const taskContext = parseTaskContext(task); - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment); - const claimLinkState = await ensureTaskClaimLink(task); + const now = nowIso() + const actor = normalizeActor(options.actor) + const force = options.force === true + const taskContext = parseTaskContext(task) + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) + const claimLinkState = await ensureTaskClaimLink(task) if ( !force && - flow.binding.prepareStatus === "ready" && + flow.binding.prepareStatus === 'ready' && flow.binding.vnId > 0 && flow.binding.vnPhone && flow.binding.bindUrl ) { if (!isKuaishouCloudBindUrlFresh(flow)) { return refreshKuaishouCloudTaskBindUrl(task, { - source: options.source || "system_refresh_expired_bind_url", + source: options.source || 'system_refresh_expired_bind_url', actor, claimLinkState, - }); + }) } - let readyTask = task; + let readyTask = task if (normalizeTaskStatus(task.task_status) !== TASK_STATUS.WAITING_BINDING) { - readyTask = (await updateTask(task.id, { - task_status: TASK_STATUS.WAITING_BINDING, - claim_token: claimLinkState.token || task.claim_token || "", - claim_expires_at: - claimLinkState.expiredAt || getTaskClaimExpiresAt(task), - updated_at: now, - })) || task; + readyTask = + (await updateTask(task.id, { + task_status: TASK_STATUS.WAITING_BINDING, + claim_token: claimLinkState.token || task.claim_token || '', + claim_expires_at: claimLinkState.expiredAt || getTaskClaimExpiresAt(task), + updated_at: now, + })) || task } return { @@ -139,222 +130,221 @@ export async function prepareKuaishouCloudFulfillmentTask( claimUrl: claimLinkState.claimUrl, token: claimLinkState.token, flow, - }; + } } - const cloudContext = resolvePersistedCloudtentaclesContextBySourceKeys( - flow.binding.cloudSourceKeys - ); - const [knapsack, skuList] = await Promise.all([ - getCloudtentaclesKnapsack(cloudContext), - listCloudtentaclesSku(cloudContext), - ]); - const resolvedBinding = resolveKuaishouCloudBindingResources(flow, { - skuItems: Array.isArray(skuList.items) ? skuList.items : [], - knapsackItems: Array.isArray(knapsack.items) ? knapsack.items : [], - }); - const vnKeyCandidates = resolveKuaishouCloudVnKeyCandidates({ - flow, - binding: resolvedBinding, - }); + const selectedCloud = await selectCloudtentaclesSourceForFulfillment(flow) + try { + const cloudContext = selectedCloud.context + const [knapsack, skuList] = await Promise.all([ + getCloudtentaclesKnapsack(cloudContext), + listCloudtentaclesSku(cloudContext), + ]) + const resolvedBinding = resolveKuaishouCloudBindingResources(flow, { + skuItems: Array.isArray(skuList.items) ? skuList.items : [], + knapsackItems: Array.isArray(knapsack.items) ? knapsack.items : [], + }) + const vnKeyCandidates = resolveKuaishouCloudVnKeyCandidates({ + flow, + binding: resolvedBinding, + }) - if (!resolvedBinding.skuId || vnKeyCandidates.length === 0) { - throw createHttpError("当前任务缺少可用 cloud 资源,无法自动准备绑定", { - statusCode: 409, - errorCode: "kuaishou_cloud_missing_binding_config", - }); - } + if (!resolvedBinding.skuId || vnKeyCandidates.length === 0) { + throw createHttpError('当前任务缺少可用 cloud 资源,无法自动准备绑定', { + statusCode: 409, + errorCode: 'kuaishou_cloud_missing_binding_config', + }) + } - const flowWithResolvedBinding = { - ...flow, - binding: { - ...flow.binding, - skuId: resolvedBinding.skuId, - skuName: resolvedBinding.skuName, - vnKey: KUAISHOU_CLOUD_FIXED_VN_KEY, - }, - }; - - const deliveryPlan = resolveKuaishouCloudDeliveryPlan(flowWithResolvedBinding, { - skuItems: Array.isArray(skuList.items) ? skuList.items : [], - knapsackItems: Array.isArray(knapsack.items) ? knapsack.items : [], - }); - if (deliveryPlan.items.length === 0) { - throw createHttpError("当前任务缺少 cloud 发货物品配置", { - statusCode: 409, - errorCode: "kuaishou_cloud_missing_delivery_items", - }); - } - - const usedKnapsack = deliveryPlan.items.every((item) => item.missingCount <= 0); - if (!usedKnapsack && !flowWithResolvedBinding.purchase.autoBuyEnabled) { - throw createHttpError("背包中没有现成库存,且当前配置未开启自动购买", { - statusCode: 409, - errorCode: "kuaishou_cloud_auto_buy_disabled", - }); - } - - const missingSku = deliveryPlan.items.find( - (item) => item.missingCount > 0 && !item.skuItem - ); - if (missingSku) { - throw createHttpError(`cloudtentacles 未找到 SKU ${missingSku.cloudSkuId}`, { - statusCode: 404, - errorCode: "kuaishou_cloud_sku_not_found", - }); - } - - const purchaseTriggered = false; - const assetBefore = 0; - const assetAfter = 0; - - const preparedBinding = await prepareKuaishouCloudBindResourceWithFallback({ - cloudContext, - vnKeyCandidates, - }); - - const nextContext = { - ...taskContext, - kuaishouCloudFulfillment: { - ...flowWithResolvedBinding, + const flowWithResolvedBinding = { + ...flow, binding: { - ...flowWithResolvedBinding.binding, - resolvedSourceKey: cloudContext.resolvedSourceKey, + ...flow.binding, + skuId: resolvedBinding.skuId, + skuName: resolvedBinding.skuName, + vnKey: KUAISHOU_CLOUD_FIXED_VN_KEY, + }, + } + + const deliveryPlan = resolveKuaishouCloudDeliveryPlan(flowWithResolvedBinding, { + skuItems: Array.isArray(skuList.items) ? skuList.items : [], + knapsackItems: Array.isArray(knapsack.items) ? knapsack.items : [], + }) + if (deliveryPlan.items.length === 0) { + throw createHttpError('当前任务缺少 cloud 发货物品配置', { + statusCode: 409, + errorCode: 'kuaishou_cloud_missing_delivery_items', + }) + } + + const usedKnapsack = deliveryPlan.items.every((item) => item.missingCount <= 0) + if (!usedKnapsack && !flowWithResolvedBinding.purchase.autoBuyEnabled) { + throw createHttpError('背包中没有现成库存,且当前配置未开启自动购买', { + statusCode: 409, + errorCode: 'kuaishou_cloud_auto_buy_disabled', + }) + } + + const missingSku = deliveryPlan.items.find((item) => item.missingCount > 0 && !item.skuItem) + if (missingSku) { + throw createHttpError(`cloudtentacles 未找到 SKU ${missingSku.cloudSkuId}`, { + statusCode: 404, + errorCode: 'kuaishou_cloud_sku_not_found', + }) + } + + const purchaseTriggered = false + const assetBefore = 0 + const assetAfter = 0 + + const preparedBinding = await prepareKuaishouCloudBindResourceWithFallback({ + cloudContext, + vnKeyCandidates, + }) + + const nextContext = { + ...taskContext, + kuaishouCloudFulfillment: { + ...flowWithResolvedBinding, + binding: { + ...flowWithResolvedBinding.binding, + resolvedSourceKey: cloudContext.resolvedSourceKey, + vnKey: preparedBinding.vnKey, + prepareStatus: 'ready', + vnId: preparedBinding.vnId, + vnPhone: preparedBinding.vnPhone, + bindUrl: preparedBinding.bindUrl, + bindPreparedAt: now, + bindExpiresAt: resolveKuaishouCloudBindUrlExpiresAt(now), + bindProbeAt: null as null, + bindProbeStatus: 'pending', + bindProbeMessage: '', + roleName: '', + roleId: '', + }, + role: { + status: 'pending', + name: '', + rid: '', + refreshedAt: null as null, + errorMessage: '', + rawInfo: null as null, + }, + purchase: { + ...flowWithResolvedBinding.purchase, + usedKnapsack, + purchaseTriggered, + assetBefore, + assetAfter, + items: deliveryPlan.items.map((item) => ({ + cloudSkuId: item.cloudSkuId, + cloudSkuName: item.cloudSkuName, + requiredCount: item.quantity, + knapsackCount: item.knapsackCount, + purchasedCount: 0, + })), + purchaseAt: flowWithResolvedBinding.purchase.purchaseAt, + }, + }, + } + + const updatedTask = await updateTask(task.id, { + task_status: TASK_STATUS.WAITING_BINDING, + user_action_status: 'pending_claim', + claim_token: claimLinkState.token || task.claim_token || '', + claim_expires_at: claimLinkState.expiredAt || getTaskClaimExpiresAt(task), + role_id: '', + role_name: '', + last_error: '', + context_json: JSON.stringify(nextContext), + updated_at: now, + }) + + await createTaskEvent( + task.id, + 'kuaishou_cloud_binding_prepared', + { + source: String(options.source || 'system').trim() || 'system', + skuId: flowWithResolvedBinding.binding.skuId, + skuName: flowWithResolvedBinding.binding.skuName, vnKey: preparedBinding.vnKey, - prepareStatus: "ready", vnId: preparedBinding.vnId, - vnPhone: preparedBinding.vnPhone, - bindUrl: preparedBinding.bindUrl, - bindPreparedAt: now, - bindExpiresAt: resolveKuaishouCloudBindUrlExpiresAt(now), - bindProbeAt: null as null, - bindProbeStatus: "pending", - bindProbeMessage: "", - roleName: "", - roleId: "", - }, - role: { - status: "pending", - name: "", - rid: "", - refreshedAt: null as null, - errorMessage: "", - rawInfo: null as null, - }, - purchase: { - ...flowWithResolvedBinding.purchase, - usedKnapsack, + vnPhoneMasked: maskPhone(preparedBinding.vnPhone), purchaseTriggered, - assetBefore, - assetAfter, - items: deliveryPlan.items.map((item) => ({ + usedKnapsack, + deliveryItems: deliveryPlan.items.map((item) => ({ cloudSkuId: item.cloudSkuId, cloudSkuName: item.cloudSkuName, - requiredCount: item.quantity, + quantity: item.quantity, knapsackCount: item.knapsackCount, purchasedCount: 0, })), - purchaseAt: flowWithResolvedBinding.purchase.purchaseAt, + resolvedByName: resolvedBinding.resolvedByName, + accountSelection: buildCloudtentaclesSourceSelectionEventDetail(selectedCloud), + actor, }, - }, - }; + now, + ) - const updatedTask = await updateTask(task.id, { - task_status: TASK_STATUS.WAITING_BINDING, - user_action_status: "pending_claim", - claim_token: claimLinkState.token || task.claim_token || "", - claim_expires_at: claimLinkState.expiredAt || getTaskClaimExpiresAt(task), - role_id: "", - role_name: "", - last_error: "", - context_json: JSON.stringify(nextContext), - updated_at: now, - }); - - await createTaskEvent( - task.id, - "kuaishou_cloud_binding_prepared", - { - source: String(options.source || "system").trim() || "system", - skuId: flowWithResolvedBinding.binding.skuId, - skuName: flowWithResolvedBinding.binding.skuName, - vnKey: preparedBinding.vnKey, - vnId: preparedBinding.vnId, - vnPhoneMasked: maskPhone(preparedBinding.vnPhone), - purchaseTriggered, - usedKnapsack, - deliveryItems: deliveryPlan.items.map((item) => ({ - cloudSkuId: item.cloudSkuId, - cloudSkuName: item.cloudSkuName, - quantity: item.quantity, - knapsackCount: item.knapsackCount, - purchasedCount: 0, - })), - resolvedByName: resolvedBinding.resolvedByName, - actor, - }, - now - ); - - return { - task: updatedTask, - claimUrl: claimLinkState.claimUrl, - token: claimLinkState.token, - flow: normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment), - }; + return { + task: updatedTask, + claimUrl: claimLinkState.claimUrl, + token: claimLinkState.token, + flow: normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment), + } + } finally { + selectedCloud.release() + } } -export async function refreshKuaishouCloudTaskBindUrl( - task: TaskRow, - options: JsonObject = {} -) { +export async function refreshKuaishouCloudTaskBindUrl(task: TaskRow, options: JsonObject = {}) { if (!isKuaishouCloudTask(task)) { - throw createHttpError("当前任务不是快手 Cloud 履约任务", { + throw createHttpError('当前任务不是快手 Cloud 履约任务', { statusCode: 409, - errorCode: "kuaishou_cloud_task_invalid", - }); + errorCode: 'kuaishou_cloud_task_invalid', + }) } - const now = nowIso(); - const actor = normalizeActor(options.actor); - const taskContext = parseTaskContext(task); - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment); + const now = nowIso() + const actor = normalizeActor(options.actor) + const taskContext = parseTaskContext(task) + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) if (!flow.binding.vnId || !flow.binding.vnKey) { - throw createHttpError("当前任务缺少可刷新绑定链接的虚拟号信息", { + throw createHttpError('当前任务缺少可刷新绑定链接的虚拟号信息', { statusCode: 409, - errorCode: "kuaishou_cloud_missing_bind_url_context", - }); + errorCode: 'kuaishou_cloud_missing_bind_url_context', + }) } const cloudContext = resolvePersistedCloudtentaclesContextBySourceKeys([ flow.binding.resolvedSourceKey, ...flow.binding.cloudSourceKeys, - ]); - const oldVnKey = flow.binding.vnKey; - const oldVnId = flow.binding.vnId; - const oldVnPhone = flow.binding.vnPhone; + ]) + const oldVnKey = flow.binding.vnKey + const oldVnId = flow.binding.vnId + const oldVnPhone = flow.binding.vnPhone await backCloudtentaclesVirtualNumber({ ...cloudContext, key: oldVnKey, id: oldVnId, - }); + }) await createTaskEvent( task.id, - "kuaishou_cloud_expired_bind_number_returned", + 'kuaishou_cloud_expired_bind_number_returned', { - source: String(options.source || "system").trim() || "system", + source: String(options.source || 'system').trim() || 'system', vnKey: oldVnKey, vnId: oldVnId, vnPhoneMasked: maskPhone(oldVnPhone), actor, }, - now - ); + now, + ) - let preparedBinding; + let preparedBinding try { preparedBinding = await prepareKuaishouCloudBindResourceWithFallback({ cloudContext, @@ -362,13 +352,13 @@ export async function refreshKuaishouCloudTaskBindUrl( flow, binding: flow.binding, }), - }); + }) - if (!String(preparedBinding.bindUrl || "").trim()) { - throw createHttpError("cloudtentacles 未返回新的绑定链接", { + if (!String(preparedBinding.bindUrl || '').trim()) { + throw createHttpError('cloudtentacles 未返回新的绑定链接', { statusCode: 502, - errorCode: "kuaishou_cloud_empty_bind_url", - }); + errorCode: 'kuaishou_cloud_empty_bind_url', + }) } } catch (error) { const failedTask = await markKuaishouCloudBindUrlRefreshFailed(task, { @@ -377,18 +367,14 @@ export async function refreshKuaishouCloudTaskBindUrl( now, actor, error, - }); + }) return { task: failedTask, - claimUrl: buildClaimUrl( - String(task.primary_claim_token || task.claim_token || "") - ), - token: String(task.primary_claim_token || task.claim_token || ""), - flow: normalizeKuaishouCloudFlow( - parseTaskContext(failedTask).kuaishouCloudFulfillment - ), - }; + claimUrl: buildClaimUrl(String(task.primary_claim_token || task.claim_token || '')), + token: String(task.primary_claim_token || task.claim_token || ''), + flow: normalizeKuaishouCloudFlow(parseTaskContext(failedTask).kuaishouCloudFulfillment), + } } const nextContext = { @@ -397,7 +383,7 @@ export async function refreshKuaishouCloudTaskBindUrl( ...flow, binding: { ...flow.binding, - prepareStatus: "ready", + prepareStatus: 'ready', vnKey: preparedBinding.vnKey, vnId: preparedBinding.vnId, vnPhone: preparedBinding.vnPhone, @@ -405,58 +391,53 @@ export async function refreshKuaishouCloudTaskBindUrl( bindPreparedAt: now, bindExpiresAt: resolveKuaishouCloudBindUrlExpiresAt(now), bindProbeAt: null as null, - bindProbeStatus: "pending", - bindProbeMessage: "", + bindProbeStatus: 'pending', + bindProbeMessage: '', roleName: normalizeTaskStatus(task.task_status) === TASK_STATUS.ROLE_CONFIRMED ? flow.binding.roleName - : "", + : '', roleId: normalizeTaskStatus(task.task_status) === TASK_STATUS.ROLE_CONFIRMED ? flow.binding.roleId - : "", + : '', }, role: normalizeTaskStatus(task.task_status) === TASK_STATUS.ROLE_CONFIRMED ? flow.role : { - status: "pending", - name: "", - rid: "", + status: 'pending', + name: '', + rid: '', refreshedAt: null, - errorMessage: "", + errorMessage: '', rawInfo: null, }, }, - }; + } - const claimLinkState = - options.claimLinkState || (await ensureTaskClaimLink(task)); + const claimLinkState = options.claimLinkState || (await ensureTaskClaimLink(task)) const updatedTask = await updateTask(task.id, { task_status: normalizeTaskStatus(task.task_status) === TASK_STATUS.ROLE_CONFIRMED ? task.task_status : TASK_STATUS.WAITING_BINDING, - claim_token: claimLinkState.token || task.claim_token || "", + claim_token: claimLinkState.token || task.claim_token || '', claim_expires_at: claimLinkState.expiredAt || getTaskClaimExpiresAt(task), role_id: - normalizeTaskStatus(task.task_status) === TASK_STATUS.ROLE_CONFIRMED - ? task.role_id - : "", + normalizeTaskStatus(task.task_status) === TASK_STATUS.ROLE_CONFIRMED ? task.role_id : '', role_name: - normalizeTaskStatus(task.task_status) === TASK_STATUS.ROLE_CONFIRMED - ? task.role_name - : "", - last_error: "", + normalizeTaskStatus(task.task_status) === TASK_STATUS.ROLE_CONFIRMED ? task.role_name : '', + last_error: '', context_json: JSON.stringify(nextContext), updated_at: now, - }); + }) await createTaskEvent( task.id, - "kuaishou_cloud_bind_url_refreshed", + 'kuaishou_cloud_bind_url_refreshed', { - source: String(options.source || "system").trim() || "system", + source: String(options.source || 'system').trim() || 'system', oldVnId, oldVnPhoneMasked: maskPhone(oldVnPhone), vnKey: preparedBinding.vnKey, @@ -464,43 +445,39 @@ export async function refreshKuaishouCloudTaskBindUrl( vnPhoneMasked: maskPhone(preparedBinding.vnPhone), actor, }, - now - ); + now, + ) return { task: updatedTask, claimUrl: claimLinkState.claimUrl, token: claimLinkState.token, flow: normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment), - }; + } } -export async function probeKuaishouCloudTaskBindUrl( - task: TaskRow, - options: JsonObject = {} -) { +export async function probeKuaishouCloudTaskBindUrl(task: TaskRow, options: JsonObject = {}) { if (!isKuaishouCloudTask(task)) { - throw createHttpError("当前任务不是快手 Cloud 履约任务", { + throw createHttpError('当前任务不是快手 Cloud 履约任务', { statusCode: 409, - errorCode: "kuaishou_cloud_task_invalid", - }); + errorCode: 'kuaishou_cloud_task_invalid', + }) } - const now = nowIso(); - const taskContext = parseTaskContext(task); - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment); + const now = nowIso() + const taskContext = parseTaskContext(task) + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) if (!flow.binding.bindUrl) { return { task, flow, probe: null as null, - }; + } } const probeIntervalMs = - Number(resolveCloudtentaclesConfig().bindUrlProbeIntervalSeconds || 30) * - 1000; + Number(resolveCloudtentaclesConfig().bindUrlProbeIntervalSeconds || 30) * 1000 if ( !options.force && flow.binding.bindProbeAt && @@ -510,26 +487,26 @@ export async function probeKuaishouCloudTaskBindUrl( task, flow, probe: null as null, - }; + } } const probe = await probeCloudtentaclesBindUrl({ bindUrl: flow.binding.bindUrl, - }); + }) if (probe.expired || !isKuaishouCloudBindUrlFresh(flow)) { const refreshed = await refreshKuaishouCloudTaskBindUrl(task, { - source: options.source || "claim_page_bind_url_expired", - actor: options.actor || { source: "system" }, - }); + source: options.source || 'claim_page_bind_url_expired', + actor: options.actor || { source: 'system' }, + }) return { ...refreshed, probe, - }; + } } - const probeRoleInfo = normalizeKuaishouCloudRoleInfo(probe.roleInfo); - const hasRoleInfo = Boolean(probeRoleInfo.name || probeRoleInfo.rid); + const probeRoleInfo = normalizeKuaishouCloudRoleInfo(probe.roleInfo) + const hasRoleInfo = Boolean(probeRoleInfo.name || probeRoleInfo.rid) const nextContext = { ...taskContext, kuaishouCloudFulfillment: { @@ -537,35 +514,35 @@ export async function probeKuaishouCloudTaskBindUrl( binding: { ...flow.binding, bindProbeAt: now, - bindProbeStatus: probe.valid ? "valid" : "invalid", - bindProbeMessage: String(probe.message || probe.reason || "").trim(), + bindProbeStatus: probe.valid ? 'valid' : 'invalid', + bindProbeMessage: String(probe.message || probe.reason || '').trim(), roleName: hasRoleInfo ? probeRoleInfo.name : flow.binding.roleName, roleId: hasRoleInfo ? probeRoleInfo.rid : flow.binding.roleId, }, role: { ...flow.role, - status: hasRoleInfo ? "ready" : flow.role.status, + status: hasRoleInfo ? 'ready' : flow.role.status, name: hasRoleInfo ? probeRoleInfo.name : flow.role.name, rid: hasRoleInfo ? probeRoleInfo.rid : flow.role.rid, refreshedAt: hasRoleInfo ? now : flow.role.refreshedAt, - errorMessage: hasRoleInfo ? "" : flow.role.errorMessage, + errorMessage: hasRoleInfo ? '' : flow.role.errorMessage, rawInfo: hasRoleInfo ? probeRoleInfo.rawInfo : flow.role.rawInfo, }, }, - }; + } const updatedTask = await updateTask(task.id, { role_id: hasRoleInfo ? probeRoleInfo.rid : task.role_id, role_name: hasRoleInfo ? probeRoleInfo.name : task.role_name, context_json: JSON.stringify(nextContext), updated_at: now, - }); + }) return { task: updatedTask, flow: normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment), probe, - }; + } } /** @@ -574,111 +551,102 @@ export async function probeKuaishouCloudTaskBindUrl( */ async function markKuaishouCloudBindUrlRefreshFailed( task: TaskRow, - { taskContext, flow, now, actor, error }: JsonObject = {} + { taskContext, flow, now, actor, error }: JsonObject = {}, ) { const errorMessage = - error instanceof Error - ? error.message - : String(error || "新绑定链接准备失败"); + error instanceof Error ? error.message : String(error || '新绑定链接准备失败') const nextContext = { ...taskContext, kuaishouCloudFulfillment: { ...flow, binding: { ...flow.binding, - prepareStatus: "pending", + prepareStatus: 'pending', vnId: 0, - vnPhone: "", - bindUrl: "", + vnPhone: '', + bindUrl: '', bindPreparedAt: null, bindExpiresAt: null, bindProbeAt: now, - bindProbeStatus: "refresh_failed", + bindProbeStatus: 'refresh_failed', bindProbeMessage: errorMessage, - roleName: "", - roleId: "", + roleName: '', + roleId: '', }, role: { - status: "pending", - name: "", - rid: "", + status: 'pending', + name: '', + rid: '', refreshedAt: now, - errorMessage: - "绑定链接已过期,旧号码已退还,新链接准备失败,请稍后刷新或联系客服处理", + errorMessage: '绑定链接已过期,旧号码已退还,新链接准备失败,请稍后刷新或联系客服处理', rawInfo: null, }, }, - }; + } const updatedTask = await updateTask(task.id, { task_status: TASK_STATUS.PENDING_BINDING_PREPARE, - role_id: "", - role_name: "", + role_id: '', + role_name: '', last_error: `绑定链接过期,旧号码已退还,新链接准备失败:${errorMessage}`, context_json: JSON.stringify(nextContext), updated_at: now, - }); + }) await createTaskEvent( task.id, - "kuaishou_cloud_bind_url_refresh_failed", + 'kuaishou_cloud_bind_url_refresh_failed', { - source: "system_refresh_expired_bind_url", + source: 'system_refresh_expired_bind_url', errorMessage, actor, }, - now - ); + now, + ) await notifyKuaishouCloudBindUrlRefreshFailed({ task: updatedTask, errorMessage, - }); + }) - return updatedTask; + return updatedTask } -export async function refreshKuaishouCloudTaskRoleInfo( - task: TaskRow, - options: JsonObject = {} -) { +export async function refreshKuaishouCloudTaskRoleInfo(task: TaskRow, options: JsonObject = {}) { if (!isKuaishouCloudTask(task)) { - throw createHttpError("当前任务不是快手 Cloud 履约任务", { + throw createHttpError('当前任务不是快手 Cloud 履约任务', { statusCode: 409, - errorCode: "kuaishou_cloud_task_invalid", - }); + errorCode: 'kuaishou_cloud_task_invalid', + }) } - const now = nowIso(); - const actor = normalizeActor(options.actor); - const taskContext = parseTaskContext(task); - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment); + const now = nowIso() + const actor = normalizeActor(options.actor) + const taskContext = parseTaskContext(task) + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) if (!flow.binding.vnId || !flow.binding.vnKey) { - throw createHttpError( - "当前任务还没有可查询的绑定角色信息,请先准备绑定资源", - { - statusCode: 409, - errorCode: "kuaishou_cloud_missing_bind_info_context", - } - ); + throw createHttpError('当前任务还没有可查询的绑定角色信息,请先准备绑定资源', { + statusCode: 409, + errorCode: 'kuaishou_cloud_missing_bind_info_context', + }) } if (flow.binding.bindUrl) { const probed = await probeKuaishouCloudTaskBindUrl(task, { - source: options.source || "system_role_refresh_bind_url_probe", + source: options.source || 'system_role_refresh_bind_url_probe', actor, force: options.forceProbe === true, - }); + }) const probedFlow = normalizeKuaishouCloudFlow( - parseTaskContext(probed.task).kuaishouCloudFulfillment - ); + parseTaskContext(probed.task).kuaishouCloudFulfillment, + ) if (probed.probe?.expired) { return { task: probed.task, roleInfo: normalizeKuaishouCloudRoleInfo(null), flow: probedFlow, - }; + } } if ( @@ -693,116 +661,98 @@ export async function refreshKuaishouCloudTaskRoleInfo( rawInfo: probedFlow.role.rawInfo, }, flow: probedFlow, - }; + } } } const cloudContext = resolvePersistedCloudtentaclesContextBySourceKeys([ flow.binding.resolvedSourceKey, ...flow.binding.cloudSourceKeys, - ]); + ]) const bindInfoResult = await getCloudtentaclesBindInfo({ ...cloudContext, key: flow.binding.vnKey, id: flow.binding.vnId, - }); + }) - const bindInfo = normalizeKuaishouCloudRoleInfo(bindInfoResult.bindInfo); - const hasRoleInfo = Boolean(bindInfo.name || bindInfo.rid); + const bindInfo = normalizeKuaishouCloudRoleInfo(bindInfoResult.bindInfo) + const hasRoleInfo = Boolean(bindInfo.name || bindInfo.rid) const nextContext = { ...taskContext, kuaishouCloudFulfillment: { ...flow, binding: { ...flow.binding, - roleName: hasRoleInfo ? bindInfo.name : "", - roleId: hasRoleInfo ? bindInfo.rid : "", + roleName: hasRoleInfo ? bindInfo.name : '', + roleId: hasRoleInfo ? bindInfo.rid : '', }, role: { - status: hasRoleInfo ? "ready" : "pending", - name: hasRoleInfo ? bindInfo.name : "", - rid: hasRoleInfo ? bindInfo.rid : "", + status: hasRoleInfo ? 'ready' : 'pending', + name: hasRoleInfo ? bindInfo.name : '', + rid: hasRoleInfo ? bindInfo.rid : '', refreshedAt: now, - errorMessage: hasRoleInfo - ? "" - : "当前还没有查询到角色信息,请完成绑定后稍等片刻再试", + errorMessage: hasRoleInfo ? '' : '当前还没有查询到角色信息,请完成绑定后稍等片刻再试', rawInfo: bindInfo.rawInfo, }, }, - }; + } const updatedTask = await updateTask(task.id, { - role_id: hasRoleInfo ? bindInfo.rid : "", - role_name: hasRoleInfo ? bindInfo.name : "", + role_id: hasRoleInfo ? bindInfo.rid : '', + role_name: hasRoleInfo ? bindInfo.name : '', context_json: JSON.stringify(nextContext), updated_at: now, - }); + }) if (options.recordEvent !== false) { await createTaskEvent( task.id, - "kuaishou_cloud_role_info_refreshed", + 'kuaishou_cloud_role_info_refreshed', { - source: String(options.source || "system").trim() || "system", + source: String(options.source || 'system').trim() || 'system', roleName: bindInfo.name, roleId: bindInfo.rid, vnId: flow.binding.vnId, actor, }, - now - ); + now, + ) } return { task: updatedTask, roleInfo: bindInfo, flow: normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment), - }; + } } function resolveKuaishouCloudDeliveryPlan( flow: JsonObject, - { - skuItems = [], - knapsackItems = [], - }: { skuItems?: unknown[]; knapsackItems?: unknown[] } = {} + { skuItems = [], knapsackItems = [] }: { skuItems?: unknown[]; knapsackItems?: unknown[] } = {}, ): { items: CloudDeliveryPlanItem[] } { - const normalizedSkuItems = Array.isArray(skuItems) - ? skuItems.filter(isCloudSkuLikeItem) - : []; + const normalizedSkuItems = Array.isArray(skuItems) ? skuItems.filter(isCloudSkuLikeItem) : [] const normalizedKnapsackItems = Array.isArray(knapsackItems) ? knapsackItems.filter(isCloudSkuLikeItem) - : []; - const deliveryItems = Array.isArray(flow.deliveryItems) - ? flow.deliveryItems - : []; + : [] + const deliveryItems = Array.isArray(flow.deliveryItems) ? flow.deliveryItems : [] return { items: deliveryItems .map((item: unknown) => { - const source = item && typeof item === "object" ? item as JsonObject : {}; - const cloudSkuId = Number(source.cloudSkuId || source.skuId || 0) || 0; - const quantity = Number(source.quantity || 1) || 1; + const source = item && typeof item === 'object' ? (item as JsonObject) : {} + const cloudSkuId = Number(source.cloudSkuId || source.skuId || 0) || 0 + const quantity = Number(source.quantity || 1) || 1 if (!cloudSkuId || !Number.isInteger(quantity) || quantity <= 0) { - return null; + return null } - const skuItem = - normalizedSkuItems.find( - (sku) => Number(sku.id || 0) === cloudSkuId - ) || null; + const skuItem = normalizedSkuItems.find((sku) => Number(sku.id || 0) === cloudSkuId) || null const knapsackItem = - normalizedKnapsackItems.find( - (sku) => Number(sku.id || 0) === cloudSkuId - ) || null; - const knapsackCount = Math.max(0, Number(knapsackItem?.count || 0) || 0); + normalizedKnapsackItems.find((sku) => Number(sku.id || 0) === cloudSkuId) || null + const knapsackCount = Math.max(0, Number(knapsackItem?.count || 0) || 0) const cloudSkuName = String( - source.cloudSkuName || - source.skuName || - skuItem?.name || - knapsackItem?.name || - "" - ).trim(); + source.cloudSkuName || source.skuName || skuItem?.name || knapsackItem?.name || '', + ).trim() return { cloudSkuId, @@ -812,18 +762,18 @@ function resolveKuaishouCloudDeliveryPlan( knapsackItem, knapsackCount, missingCount: Math.max(0, quantity - knapsackCount), - }; + } }) .filter((item): item is CloudDeliveryPlanItem => Boolean(item)), - }; + } } function isCloudSkuLikeItem(item: unknown): item is JsonObject { - const current = item && typeof item === "object" ? item as JsonObject : {}; - return Number(current.id || 0) > 0; + const current = item && typeof item === 'object' ? (item as JsonObject) : {} + return Number(current.id || 0) > 0 } export { dispatchKuaishouCloudFulfillmentTask, returnKuaishouCloudFulfillmentTask, -} from "./task-finalization.js"; +} from './task-finalization.js' diff --git a/apps/backend/src/services/order/delivery-task-service.test.ts b/apps/backend/src/services/order/delivery-task-service.test.ts new file mode 100644 index 00000000..ea56e3db --- /dev/null +++ b/apps/backend/src/services/order/delivery-task-service.test.ts @@ -0,0 +1,129 @@ +import assert from 'node:assert/strict' +import test from 'node:test' + +import { syncDeliveryTasksForOrderWithDeps } from './delivery-task-service.js' + +test('syncDeliveryTasksForOrder 多账号候选时不提前固定 cloudtentacles 账号', async () => { + const createdInputs: any[] = [] + + await syncDeliveryTasksForOrderWithDeps( + { + id: 1, + pay_status: 'pending_payment', + provider: 'open_91', + platform: 'kuaishou', + shop_id: 'shop-1', + shop_name: '测试店铺', + platform_order_id: '2614900602069169', + } as any, + [ + { + id: 10, + order_id: 1, + sku_code: 'SKU-1', + sku_name: '荣耀勋章礼包(30个)', + quantity: 1, + item_snapshot_json: { + cloudtentacles: { + cloudSourceKeys: ['account-a', 'account-b'], + resolvedSourceKey: 'account-a', + cloudSkuId: 74, + cloudSkuName: '荣耀勋章礼包(30个)', + deliveryItems: [ + { + cloudSkuId: 74, + cloudSkuName: '荣耀勋章礼包(30个)', + quantity: 1, + }, + ], + }, + }, + }, + ] as any, + { + listTasksByOrderId: async () => [], + getFulfillmentProfileByKey: async () => ({ + id: 9, + profile_key: 'kuaishou_ct_assisted', + profile_name: '快手 cloud 履约', + executor_key: 'kuaishou_ct_assisted', + }), + createTask: async (input: any) => { + createdInputs.push(input) + return { + id: 100, + order_item_id: input.orderItemId, + context_json: input.contextJson, + task_status: input.taskStatus, + executor_key: input.executorKey, + } as any + }, + nowIso: () => '2026-05-29T09:05:46.658Z', + randomId: () => 'DT-test', + }, + ) + + const context = JSON.parse(createdInputs[0].contextJson) + assert.deepEqual(context.kuaishouCloudFulfillment.binding.cloudSourceKeys, [ + 'account-a', + 'account-b', + ]) + assert.equal(context.kuaishouCloudFulfillment.binding.resolvedSourceKey, '') +}) + +test('syncDeliveryTasksForOrder 单账号候选时保留固定 cloudtentacles 账号', async () => { + const createdInputs: any[] = [] + + await syncDeliveryTasksForOrderWithDeps( + { + id: 2, + pay_status: 'pending_payment', + provider: 'open_91', + platform: 'kuaishou', + shop_id: 'shop-1', + shop_name: '测试店铺', + platform_order_id: '2614900602069170', + } as any, + [ + { + id: 11, + order_id: 2, + sku_code: 'SKU-1', + sku_name: '荣耀勋章礼包(30个)', + quantity: 1, + item_snapshot_json: { + cloudtentacles: { + cloudSourceKeys: ['account-a'], + cloudSkuId: 74, + cloudSkuName: '荣耀勋章礼包(30个)', + }, + }, + }, + ] as any, + { + listTasksByOrderId: async () => [], + getFulfillmentProfileByKey: async () => ({ + id: 9, + profile_key: 'kuaishou_ct_assisted', + profile_name: '快手 cloud 履约', + executor_key: 'kuaishou_ct_assisted', + }), + createTask: async (input: any) => { + createdInputs.push(input) + return { + id: 101, + order_item_id: input.orderItemId, + context_json: input.contextJson, + task_status: input.taskStatus, + executor_key: input.executorKey, + } as any + }, + nowIso: () => '2026-05-29T09:05:46.658Z', + randomId: () => 'DT-test', + }, + ) + + const context = JSON.parse(createdInputs[0].contextJson) + assert.deepEqual(context.kuaishouCloudFulfillment.binding.cloudSourceKeys, ['account-a']) + assert.equal(context.kuaishouCloudFulfillment.binding.resolvedSourceKey, 'account-a') +}) diff --git a/apps/backend/src/services/order/delivery-task-service.ts b/apps/backend/src/services/order/delivery-task-service.ts index cf9f717f..297f98d3 100644 --- a/apps/backend/src/services/order/delivery-task-service.ts +++ b/apps/backend/src/services/order/delivery-task-service.ts @@ -53,13 +53,12 @@ type DeliveryTaskDeps = { randomId?: (prefix?: string) => string } -type RuntimeDeliveryTaskDeps = Required> +type RuntimeDeliveryTaskDeps = Required< + Pick< + DeliveryTaskDeps, + 'updateTask' | 'createTaskClaimToken' | 'notifyTaskAutoManualReview' | 'nowIso' + > +> type TaskContext = { [key: string]: unknown @@ -103,11 +102,18 @@ export async function syncDeliveryTasksForOrderWithDeps( } const itemMap = new Map(orderItems.map((item) => [item.id, item])) - const preparedTasks = await Promise.all(existingTasks.map((task) => preparePaidTask({ - ...task, - skuCode: itemMap.get(task.order_item_id)?.sku_code || '', - skuName: itemMap.get(task.order_item_id)?.sku_name || '', - }, runtimeDeps))) + const preparedTasks = await Promise.all( + existingTasks.map((task) => + preparePaidTask( + { + ...task, + skuCode: itemMap.get(task.order_item_id)?.sku_code || '', + skuName: itemMap.get(task.order_item_id)?.sku_name || '', + }, + runtimeDeps, + ), + ), + ) return preparedTasks.filter(isTaskRow) } @@ -125,6 +131,11 @@ export async function syncDeliveryTasksForOrderWithDeps( const cloudtentaclesConfig = parseJsonObject(fulfillmentConfig.cloudtentacles) const kuaishouConsumeConfig = parseJsonObject(fulfillmentConfig.kuaishouConsume) const kuaishouShopConfig = parseJsonObject(fulfillmentConfig.kuaishouShop) + const cloudSourceKeys = normalizeStringArray(cloudtentaclesConfig.cloudSourceKeys) + const resolvedCloudSourceKey = + cloudSourceKeys.length === 1 + ? String(cloudtentaclesConfig.resolvedSourceKey || cloudSourceKeys[0] || '').trim() + : '' const deliveryItems = normalizeCloudDeliveryItems(cloudtentaclesConfig) const primaryDeliveryItem = deliveryItems[0] || { cloudSkuId: Number(cloudtentaclesConfig.skuId || 0) || 0, @@ -134,9 +145,8 @@ export async function syncDeliveryTasksForOrderWithDeps( for (let index = 0; index < quantity; index += 1) { const createdAt = getNowIso() - const initialStatus = order.pay_status === 'paid' - ? resolvePaidTaskStatus(profile) - : TASK_STATUS.PENDING_PAYMENT + const initialStatus = + order.pay_status === 'paid' ? resolvePaidTaskStatus(profile) : TASK_STATUS.PENDING_PAYMENT const task = await createDeliveryTask({ orderId: order.id, @@ -168,78 +178,90 @@ export async function syncDeliveryTasksForOrderWithDeps( skuName: item.sku_name, kuaishouCloudFulfillment: isKuaishouCloudExecutor(profile.executor_key) ? { - flowType: 'kuaishou_cloud_fulfillment', - configId: String(fulfillmentConfig.configId || '').trim(), - internalSkuCode: item.sku_code, - internalSkuName: item.sku_name, - deliveryItems, - ticket: { - code: '', - status: 'pending', - capturedAt: null, - capturedBy: null, - verifiedAt: null, - oid: '', - formToken: '', - leftCount: 0, - goodsTitle: '', - }, - binding: { - prepareStatus: 'pending', - cloudSourceKeys: normalizeStringArray(cloudtentaclesConfig.cloudSourceKeys), - resolvedSourceKey: String(cloudtentaclesConfig.resolvedSourceKey || '').trim(), - skuId: primaryDeliveryItem.cloudSkuId, - skuName: primaryDeliveryItem.cloudSkuName, - vnKey: '1', - vnId: 0, - vnPhone: '', - bindUrl: '', - bindPreparedAt: null, - bindExpiresAt: null, - bindProbeAt: null, - bindProbeStatus: '', - bindProbeMessage: '', - }, - role: { - status: 'pending', - name: '', - rid: '', - refreshedAt: null, - errorMessage: '', - rawInfo: null, - }, - purchase: { - autoBuyEnabled: cloudtentaclesConfig.autoBuyEnabled !== false, - minAssetReserve: Number(cloudtentaclesConfig.minAssetReserve || 0) || 0, - usedKnapsack: false, - purchaseTriggered: false, - assetBefore: 0, - assetAfter: 0, - purchaseAt: null, - }, - dispatch: { - status: 'pending', - dispatchAt: null, - dispatchBy: null, - sendType: 0, - note: '', - }, - returnNumber: { - status: 'pending', - returnedAt: null, - returnedBy: null, - autoReturnEnabled: cloudtentaclesConfig.autoReturnNumberAfterDispatch === true, - }, - consume: { - status: 'pending', - shopId: String(kuaishouConsumeConfig.shopId || kuaishouShopConfig.shopId || itemSnapshot.shopId || order.shop_id || '').trim(), - shopName: String(kuaishouConsumeConfig.shopName || kuaishouShopConfig.kshopName || itemSnapshot.shopName || order.shop_name || '').trim(), - autoConsumeEnabled: kuaishouConsumeConfig.autoConsumeAfterDispatch === true, - consumedAt: null, - errorMessage: '', - }, - notes: String(fulfillmentConfig.notes || '').trim(), - } + flowType: 'kuaishou_cloud_fulfillment', + configId: String(fulfillmentConfig.configId || '').trim(), + internalSkuCode: item.sku_code, + internalSkuName: item.sku_name, + deliveryItems, + ticket: { + code: '', + status: 'pending', + capturedAt: null, + capturedBy: null, + verifiedAt: null, + oid: '', + formToken: '', + leftCount: 0, + goodsTitle: '', + }, + binding: { + prepareStatus: 'pending', + cloudSourceKeys, + resolvedSourceKey: resolvedCloudSourceKey, + skuId: primaryDeliveryItem.cloudSkuId, + skuName: primaryDeliveryItem.cloudSkuName, + vnKey: '1', + vnId: 0, + vnPhone: '', + bindUrl: '', + bindPreparedAt: null, + bindExpiresAt: null, + bindProbeAt: null, + bindProbeStatus: '', + bindProbeMessage: '', + }, + role: { + status: 'pending', + name: '', + rid: '', + refreshedAt: null, + errorMessage: '', + rawInfo: null, + }, + purchase: { + autoBuyEnabled: cloudtentaclesConfig.autoBuyEnabled !== false, + minAssetReserve: Number(cloudtentaclesConfig.minAssetReserve || 0) || 0, + usedKnapsack: false, + purchaseTriggered: false, + assetBefore: 0, + assetAfter: 0, + purchaseAt: null, + }, + dispatch: { + status: 'pending', + dispatchAt: null, + dispatchBy: null, + sendType: 0, + note: '', + }, + returnNumber: { + status: 'pending', + returnedAt: null, + returnedBy: null, + autoReturnEnabled: cloudtentaclesConfig.autoReturnNumberAfterDispatch === true, + }, + consume: { + status: 'pending', + shopId: String( + kuaishouConsumeConfig.shopId || + kuaishouShopConfig.shopId || + itemSnapshot.shopId || + order.shop_id || + '', + ).trim(), + shopName: String( + kuaishouConsumeConfig.shopName || + kuaishouShopConfig.kshopName || + itemSnapshot.shopName || + order.shop_name || + '', + ).trim(), + autoConsumeEnabled: kuaishouConsumeConfig.autoConsumeAfterDispatch === true, + consumedAt: null, + errorMessage: '', + }, + notes: String(fulfillmentConfig.notes || '').trim(), + } : null, }), createdAt, @@ -344,7 +366,9 @@ function parseTaskContext(task: { context_json?: unknown } | null | undefined): try { const parsed = JSON.parse(String(value || '{}')) - return parsed && typeof parsed === 'object' && !Array.isArray(parsed) ? parsed as TaskContext : {} + return parsed && typeof parsed === 'object' && !Array.isArray(parsed) + ? (parsed as TaskContext) + : {} } catch { return {} } @@ -361,7 +385,9 @@ function parseJsonObject(value: unknown): JsonObject { try { const parsed = JSON.parse(String(value || '{}')) - return parsed && typeof parsed === 'object' && !Array.isArray(parsed) ? parsed as JsonObject : {} + return parsed && typeof parsed === 'object' && !Array.isArray(parsed) + ? (parsed as JsonObject) + : {} } catch { return {} } @@ -385,7 +411,11 @@ async function resolveDynamicCloudtentaclesProfile( quantity: 1, } - if (!primaryDeliveryItem.cloudSkuId || !primaryDeliveryItem.cloudSkuName || cloudSourceKeys.length === 0) { + if ( + !primaryDeliveryItem.cloudSkuId || + !primaryDeliveryItem.cloudSkuName || + cloudSourceKeys.length === 0 + ) { return null } @@ -409,7 +439,10 @@ async function resolveDynamicCloudtentaclesProfile( cloudSourceKeys, skuId: primaryDeliveryItem.cloudSkuId, skuName: primaryDeliveryItem.cloudSkuName, - resolvedSourceKey: String(cloudtentacles.resolvedSourceKey || '').trim(), + resolvedSourceKey: + cloudSourceKeys.length === 1 + ? String(cloudtentacles.resolvedSourceKey || cloudSourceKeys[0] || '').trim() + : '', deliveryItems, vnKey: '1', autoBuyEnabled: true, @@ -421,9 +454,10 @@ async function resolveDynamicCloudtentaclesProfile( shopName: String(snapshot.shopName || '').trim(), autoConsumeAfterDispatch: false, }, - notes: cloudtentacles.matchMode === 'cloudtentacles_override' - ? '91卡券商品名命中 cloudtentacles 覆盖规则' - : '91卡券商品名自动匹配 cloudtentacles 商品', + notes: + cloudtentacles.matchMode === 'cloudtentacles_override' + ? '91卡券商品名命中 cloudtentacles 覆盖规则' + : '91卡券商品名自动匹配 cloudtentacles 商品', }, } } @@ -433,9 +467,7 @@ function normalizeStringArray(value: unknown): string[] { return [] } - return Array.from( - new Set(value.map((item) => String(item || '').trim()).filter(Boolean)), - ) + return Array.from(new Set(value.map((item) => String(item || '').trim()).filter(Boolean))) } function normalizeCloudDeliveryItems(value: JsonObject): Array<{ @@ -446,7 +478,9 @@ function normalizeCloudDeliveryItems(value: JsonObject): Array<{ const rawItems = Array.isArray(value.deliveryItems) ? value.deliveryItems : [] const items = rawItems .map((item) => normalizeCloudDeliveryItem(item)) - .filter((item): item is { cloudSkuId: number; cloudSkuName: string; quantity: number } => Boolean(item)) + .filter((item): item is { cloudSkuId: number; cloudSkuName: string; quantity: number } => + Boolean(item), + ) if (items.length > 0) { return mergeCloudDeliveryItems(items) @@ -457,17 +491,18 @@ function normalizeCloudDeliveryItems(value: JsonObject): Array<{ return [] } - return [{ - cloudSkuId, - cloudSkuName: String(value.skuName || '').trim(), - quantity: 1, - }] + return [ + { + cloudSkuId, + cloudSkuName: String(value.skuName || '').trim(), + quantity: 1, + }, + ] } function normalizeCloudDeliveryItem(value: unknown) { - const source = value && typeof value === 'object' && !Array.isArray(value) - ? value as JsonObject - : {} + const source = + value && typeof value === 'object' && !Array.isArray(value) ? (value as JsonObject) : {} const cloudSkuId = Number(source.cloudSkuId || source.skuId || 0) || 0 if (!Number.isInteger(cloudSkuId) || cloudSkuId <= 0) { return null