重构 lewan 履约:拆分发货模块并统一 registry 入口

将 task-finalization 拆为 dispatch/stock/return/redeem/confirm 等用例,
claim 与 admin 经 executor registry 调用;admin 退号仅清理虚拟号不核销。
This commit is contained in:
yml2213
2026-07-10 22:23:26 +08:00
parent c665ef5851
commit 792680bd19
19 changed files with 2111 additions and 1488 deletions
@@ -7,10 +7,6 @@ import {
listKuaishouIndustryVouchersByOid,
listKuaishouIndustryVouchersByTaskId,
} from '../../../repositories/kuaishou-industry-voucher-repo.js'
import {
backCloudtentaclesVirtualNumber,
getCloudtentaclesBindInfo,
} from '../../platforms/cloudtentacles/virtual-number-service.js'
import { resendKuaishouIndustryVoucherSendCallback } from '../../platforms/kuaishou-industry/send-code-service.js'
import { consumeKuaishouIndustryVoucher } from '../../platforms/kuaishou-industry/voucher-service.js'
import {
@@ -22,17 +18,16 @@ import { nowIso } from '../../../utils/time.js'
import { TASK_STATUS } from '../../../domain/task-status.js'
import { createAdminViewerContext, parseTaskContext } from '../admin-read-shared-helpers.js'
import { getRequiredTask, mapTaskActionPayload } from '../admin-task-read-helpers.js'
import { maskCode, maskPhone } from '../write-helpers.js'
import {
dispatchKuaishouCloudFulfillmentTask,
prepareKuaishouCloudFulfillmentTask,
rebindKuaishouCloudTaskRole,
} from '../../fulfillment/kuaishou-cloud/index.js'
dispatchFulfillmentTask,
prepareFulfillmentBinding,
rebindFulfillmentRole,
refreshFulfillmentRole,
returnFulfillmentNumber,
} from '../../fulfillment/executors/registry.js'
import {
isKuaishouCloudTask,
normalizeKuaishouCloudFlow,
normalizeKuaishouCloudRoleInfo,
resolvePersistedCloudtentaclesContext,
} from './kuaishou-cloud-helpers.js'
import type {
@@ -66,15 +61,21 @@ export async function prepareAdminTaskKuaishouCloudFulfillment(
})
}
const result = await prepareKuaishouCloudFulfillmentTask(task, {
const result = await prepareFulfillmentBinding(task, {
source: 'admin_task_prepare',
actor: buildAdminActionActor(session),
})
if (!result?.task) {
throw createHttpError('当前任务不支持准备绑定资源', {
statusCode: 409,
errorCode: 'admin_task_prepare_not_supported',
})
}
return {
task: mapTaskActionPayload(result.task),
claimUrl: result.claimUrl,
token: result.token,
claimUrl: String(result.claimUrl || ''),
token: String(result.token || ''),
}
}
@@ -99,11 +100,17 @@ export async function dispatchAdminTaskKuaishouCloudFulfillment(
})
}
const result = await dispatchKuaishouCloudFulfillmentTask(task, {
const result = await dispatchFulfillmentTask(task, {
source: 'admin_task_dispatch',
actor: buildAdminActionActor(session),
autoFinalize: false,
})
if (!result?.task) {
throw createHttpError('当前任务不支持 executor 发货', {
statusCode: 409,
errorCode: 'admin_task_dispatch_not_supported',
})
}
return {
task: mapTaskActionPayload(result.task),
@@ -115,7 +122,6 @@ export async function refreshAdminTaskKuaishouCloudRoleInfo(
session: AdminViewerSessionInput | null = null,
): Promise<AdminTaskActionResponse> {
const task = await getRequiredTask(taskId)
const now = nowIso()
const viewerContext = createAdminViewerContext(session)
if (!viewerContext.canOperateAssistedTask) {
@@ -132,74 +138,21 @@ export async function refreshAdminTaskKuaishouCloudRoleInfo(
})
}
const taskContext = parseTaskContext(task)
const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment)
if (!flow.binding.vnId || !flow.binding.vnKey) {
throw createHttpError('当前任务还没有可查询的绑定角色信息,请先准备绑定资源', {
const result = await refreshFulfillmentRole(task, {
source: 'admin_task_refresh_role',
actor: buildAdminActionActor(session),
forceProbe: true,
recordEvent: true,
})
if (!result?.task) {
throw createHttpError('当前任务不支持刷新角色信息', {
statusCode: 409,
errorCode: 'admin_task_kuaishou_cloud_missing_bind_info_context',
errorCode: 'admin_task_refresh_not_supported',
})
}
const cloudContext = resolvePersistedCloudtentaclesContext([
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 nextContext = {
...taskContext,
kuaishouCloudFulfillment: {
...flow,
binding: {
...flow.binding,
roleName: bindInfo.name,
roleId: bindInfo.rid,
},
role: {
status: bindInfo.name || bindInfo.rid ? 'ready' : 'pending',
name: bindInfo.name,
rid: bindInfo.rid,
refreshedAt: now,
errorMessage:
bindInfo.name || bindInfo.rid ? '' : '当前还没有查询到角色信息,请让客户完成绑定后再刷新',
rawInfo: bindInfo.rawInfo,
},
},
}
const updatedTask = await updateTask(task.id, {
role_id: bindInfo.rid || '',
role_name: bindInfo.name || '',
context_json: JSON.stringify(nextContext),
updated_at: now,
})
await createTaskEvent(
task.id,
'kuaishou_cloud_role_info_refreshed',
{
roleName: bindInfo.name,
roleId: bindInfo.rid,
vnId: flow.binding.vnId,
refreshedBy: session
? {
userId: Number(session.userId || 0) || 0,
username: String(session.username || '').trim(),
role: String(session.role || '').trim(),
}
: null,
},
now,
)
return {
task: mapTaskActionPayload(updatedTask),
task: mapTaskActionPayload(result.task),
}
}
@@ -218,21 +171,21 @@ export async function rebindAdminTaskKuaishouCloudRole(
})
}
const result = await rebindKuaishouCloudTaskRole(task, {
const result = await rebindFulfillmentRole(task, {
source: 'admin_task_rebind_role',
actor: session
? {
userId: Number(session.userId || 0) || 0,
username: String(session.username || '').trim(),
role: String(session.role || '').trim(),
}
: null,
actor: buildAdminActionActor(session),
})
if (!result?.task) {
throw createHttpError('当前任务不支持换绑角色', {
statusCode: 409,
errorCode: 'admin_task_rebind_not_supported',
})
}
return {
task: mapTaskActionPayload(result.task),
claimUrl: result.claimUrl,
token: result.token,
claimUrl: String(result.claimUrl || ''),
token: String(result.token || ''),
}
}
@@ -241,7 +194,6 @@ export async function returnNumberAdminTaskKuaishouCloudFulfillment(
session: AdminViewerSessionInput | null = null,
): Promise<AdminTaskActionResponse> {
const task = await getRequiredTask(taskId)
const now = nowIso()
const viewerContext = createAdminViewerContext(session)
if (!viewerContext.canManageTaskLifecycle) {
@@ -258,108 +210,21 @@ export async function returnNumberAdminTaskKuaishouCloudFulfillment(
})
}
const taskContext = parseTaskContext(task)
const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment)
const cloudContext = resolvePersistedCloudtentaclesContext([
flow.binding.resolvedSourceKey,
...flow.binding.cloudSourceKeys,
])
if (!flow.binding.vnId || !flow.binding.vnKey) {
throw createHttpError('当前任务缺少可退还的虚拟号信息', {
// admin 退号 = 清理虚拟号占位,不核销电子凭证
const result = await returnFulfillmentNumber(task, {
source: 'admin_task_return_number',
actor: buildAdminActionActor(session),
consumeIndustryVoucher: false,
})
if (!result?.task) {
throw createHttpError('当前任务不支持退还号码', {
statusCode: 409,
errorCode: 'admin_task_kuaishou_cloud_missing_return_context',
errorCode: 'admin_task_return_not_supported',
})
}
await backCloudtentaclesVirtualNumber({
...cloudContext,
key: flow.binding.vnKey,
id: flow.binding.vnId,
})
const ticketCode = String(flow.ticket.code || '').trim()
const consumeAlreadyCompleted = flow.consume.status === 'success'
const consumeStatus = consumeAlreadyCompleted ? 'success' : 'skipped'
let consumeErrorMessage = ''
const consumedAt = consumeAlreadyCompleted ? flow.consume.consumedAt || now : null
const nextTaskStatus = 'completed'
const nextResultCode = consumeAlreadyCompleted
? 'kuaishou_cloud_completed'
: 'kuaishou_cloud_completed_without_eticket_consume'
const nextResultMessage = consumeAlreadyCompleted
? 'cloudtentacles 发货、退号并完成电子凭证核销'
: 'cloudtentacles 发货、退号并收口'
const nextContext = {
...taskContext,
kuaishouCloudFulfillment: {
...flow,
returnNumber: {
...flow.returnNumber,
status: 'success',
returnedAt: now,
returnedBy: session
? {
userId: Number(session.userId || 0) || 0,
username: String(session.username || '').trim(),
role: String(session.role || '').trim(),
}
: null,
},
consume: {
...flow.consume,
status: consumeStatus,
shopId: flow.consume.shopId,
shopName: flow.consume.shopName,
autoConsumeEnabled: flow.consume.autoConsumeEnabled === true,
consumedAt,
errorMessage: consumeErrorMessage,
},
},
}
const updatedTask = await updateTask(task.id, {
task_status: nextTaskStatus,
delivery_status: 'delivered',
result_code: nextResultCode,
result_message: nextResultMessage,
redeemed_at: now,
last_error: consumeErrorMessage,
context_json: JSON.stringify(nextContext),
updated_at: now,
})
await createTaskEvent(
task.id,
'kuaishou_cloud_number_returned',
{
vnId: flow.binding.vnId,
vnPhoneMasked: maskPhone(flow.binding.vnPhone),
},
now,
)
const consumeEventType = consumeAlreadyCompleted
? 'kuaishou_cloud_consume_already_completed'
: 'kuaishou_cloud_consume_skipped'
await createTaskEvent(
task.id,
consumeEventType,
{
ticketCodeMasked: maskCode(ticketCode),
shopId: flow.consume.shopId,
shopName: flow.consume.shopName,
consumeStatus,
consumeMode: 'legacy_writeoff_disabled',
errorMessage: consumeErrorMessage,
},
now,
)
return {
task: mapTaskActionPayload(updatedTask),
task: mapTaskActionPayload(result.task),
}
}
@@ -1,27 +1,22 @@
import { createTaskEvent } from '../../repositories/task-event-repo.js'
import { getTaskById, updateTask, updateTaskStatusIfCurrent } from '../../repositories/task-repo.js'
import { updateTask } from '../../repositories/task-repo.js'
import { createHttpError } from '../../utils/http.js'
import { parseTaskContext as parseTaskContextValue } from '../../utils/task-json.js'
import { addHours, nowIso } from '../../utils/time.js'
import { asJsonObject, type JsonObject } from '../../types/json.js'
import { nowIso } from '../../utils/time.js'
import type { JsonObject } from '../../types/json.js'
import { TASK_STATUS } from '../../domain/task-status.js'
import {
TASK_STATUS,
canRedeemKuaishouCloudClaimStatus,
isKuaishouCloudRedeemSettledStatus,
isKuaishouCloudRoleConfirmSettledStatus,
normalizeTaskStatus,
} from '../../domain/task-status.js'
confirmFulfillmentRole,
prepareFulfillmentBinding,
rebindFulfillmentRole,
redeemFulfillmentTask,
} from '../fulfillment/executors/registry.js'
import {
dispatchKuaishouCloudFulfillmentTask,
isKuaishouCloudMockTask,
normalizeKuaishouCloudFlow,
prepareKuaishouCloudFulfillmentTask,
rebindKuaishouCloudTaskRole,
refreshKuaishouCloudTaskRoleInfo,
} from '../fulfillment/kuaishou-cloud/index.js'
import { syncKuaishouFeifeiTaskStatus } from '../fulfillment/kuaishou-feifei/index.js'
import {
assertBoundUidMatchesExpected,
assertClaimExpectedUidReady,
assertValidClaimUid,
canUpdateClaimUid,
getClaimIdentityFromTask,
@@ -32,6 +27,30 @@ import type { TaskRow } from '../../types/repository/rows.js'
type ClaimDetailPayload = ReturnType<typeof buildClaimDetailPayload>
function assertLewanClaimTask(task: TaskRow) {
if (String(task.executor_key || '').trim() !== 'kuaishou_ct_assisted') {
throw createHttpError('当前领取链接不是 kuaishou-lewan 客户领取流程', {
statusCode: 409,
errorCode: 'claim_not_kuaishou_cloud',
})
}
}
async function requireExecutorAction<T>(
result: Promise<T | null> | T | null,
errorMessage: string,
errorCode: string,
): Promise<NonNullable<T>> {
const value = await result
if (!value) {
throw createHttpError(errorMessage, {
statusCode: 409,
errorCode,
})
}
return value
}
async function verifyIndustryVoucherTicket(
context: Awaited<ReturnType<typeof getClaimContext>>,
now: string,
@@ -126,7 +145,7 @@ async function verifyIndustryVoucherTicket(
const preparedFlow = normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment)
if (preparedFlow.binding.prepareStatus !== 'ready' || !preparedFlow.binding.bindUrl) {
await prepareKuaishouCloudFulfillmentTask(taskWithTicket, {
await prepareFulfillmentBinding(taskWithTicket, {
source: 'send_code_callback',
actor: { source: 'send_code_callback' },
})
@@ -262,7 +281,7 @@ export async function submitClaimUid(
const executorKey = String(updatedTask.executor_key || '').trim()
if (executorKey === 'kuaishou_ct_assisted' && !isKuaishouCloudMockTask(updatedTask)) {
try {
const prepared = await prepareKuaishouCloudFulfillmentTask(updatedTask, {
const prepared = await prepareFulfillmentBinding(updatedTask, {
source: 'claim_page_submit_uid',
actor: { source: 'claim_page' },
})
@@ -279,18 +298,55 @@ export async function submitClaimUid(
export async function rebindKuaishouCloudClaimRole(token: unknown) {
const context = await getClaimContext(token)
assertLewanClaimTask(context.task)
if (String(context.task.executor_key || '').trim() !== 'kuaishou_ct_assisted') {
throw createHttpError('当前领取链接不是 kuaishou-lewan 客户领取流程', {
statusCode: 409,
errorCode: 'claim_not_kuaishou_cloud',
})
}
await requireExecutorAction(
rebindFulfillmentRole(context.task, {
source: 'claim_page_rebind_role',
actor: { source: 'claim_page' },
}),
'当前领取链接不支持换绑角色',
'claim_rebind_not_supported',
)
await rebindKuaishouCloudTaskRole(context.task, {
source: 'claim_page_rebind_role',
actor: { source: 'claim_page' },
})
return getKuaishouCloudClaimDetail(token)
}
export async function confirmKuaishouCloudClaimRole(token: unknown) {
const context = await getClaimContext(token)
assertLewanClaimTask(context.task)
await requireExecutorAction(
confirmFulfillmentRole(context.task, {
source: 'claim_page_role_confirm',
actor: { source: 'claim_page' },
errorCodePrefix: 'claim_kuaishou_cloud',
forceProbe: true,
}),
'当前领取链接不支持确认角色',
'claim_confirm_not_supported',
)
return getKuaishouCloudClaimDetail(token)
}
/**
* 领取页一键兑换:鉴权后交给 lewan 履约引擎 redeem 用例。
*/
export async function redeemKuaishouCloudClaim(token: unknown) {
const context = await getClaimContext(token)
assertLewanClaimTask(context.task)
await requireExecutorAction(
redeemFulfillmentTask(context.task, {
source: 'claim_page_redeem',
actor: { source: 'claim_page' },
autoFinalize: true,
errorCodePrefix: 'claim_kuaishou_cloud_redeem',
}),
'当前领取链接不支持一键兑换',
'claim_redeem_not_supported',
)
return getKuaishouCloudClaimDetail(token)
}
@@ -327,321 +383,3 @@ function normalizeIndustryVoucherContext(value: unknown): JsonObject {
consumedAt: source.consumedAt || null,
}
}
export async function confirmKuaishouCloudClaimRole(token: unknown) {
const context = await getClaimContext(token)
const now = nowIso()
if (String(context.task.executor_key || '').trim() !== 'kuaishou_ct_assisted') {
throw createHttpError('当前领取链接不是 kuaishou-lewan 客户领取流程', {
statusCode: 409,
errorCode: 'claim_not_kuaishou_cloud',
})
}
if (isKuaishouCloudRoleConfirmSettledStatus(context.task.task_status)) {
return getKuaishouCloudClaimDetail(token)
}
const expectedUid = assertClaimExpectedUidReady(context.task)
const contextSource = parseTaskContext(context.task)
const mockMode = isKuaishouCloudMockContext(contextSource)
const refreshed = mockMode
? { task: context.task }
: await refreshKuaishouCloudTaskRoleInfo(context.task, {
source: 'claim_page_role_confirm',
actor: { source: 'claim_page' },
recordEvent: false,
forceProbe: true,
})
const flow = normalizeKuaishouCloudFlow(parseTaskContext(refreshed.task).kuaishouCloudFulfillment)
if (!flow.binding.vnPhone || !flow.binding.roleName || !flow.binding.roleId) {
throw createHttpError('角色信息还未刷新到系统,请完成绑定后稍等片刻再试', {
statusCode: 409,
errorCode: 'claim_kuaishou_cloud_role_not_ready',
})
}
if (!mockMode) {
assertBoundUidMatchesExpected(refreshed.task, flow, {
errorCodePrefix: 'claim_kuaishou_cloud',
})
}
await updateTask(context.task.id, {
task_status: TASK_STATUS.ROLE_CONFIRMED,
role_id: flow.binding.roleId,
role_name: flow.binding.roleName,
user_action_status: TASK_STATUS.ROLE_CONFIRMED,
role_confirmed_at: now,
last_error: '',
updated_at: now,
})
await createTaskEvent(
context.task.id,
'kuaishou_cloud_role_confirmed',
{
expectedUid,
vnPhone: flow.binding.vnPhone,
roleName: flow.binding.roleName,
roleId: flow.binding.roleId,
},
now,
)
return getKuaishouCloudClaimDetail(token)
}
export async function redeemKuaishouCloudClaim(token: unknown) {
const context = await getClaimContext(token)
const now = nowIso()
if (String(context.task.executor_key || '').trim() !== 'kuaishou_ct_assisted') {
throw createHttpError('当前领取链接不是 kuaishou-lewan 客户领取流程', {
statusCode: 409,
errorCode: 'claim_not_kuaishou_cloud',
})
}
let workingTask = context.task
const currentStatus = normalizeTaskStatus(workingTask.task_status)
if (isKuaishouCloudRedeemSettledStatus(currentStatus)) {
return getKuaishouCloudClaimDetail(token)
}
if (!canRedeemKuaishouCloudClaimStatus(currentStatus)) {
throw createHttpError('当前状态不可兑换,请先完成绑定并匹配 UID', {
statusCode: 409,
errorCode: 'claim_kuaishou_cloud_not_ready_for_redeem',
})
}
// WAITING_BINDING + UID 匹配:自动升为 ROLE_CONFIRMED,实现一键兑换
if (currentStatus === TASK_STATUS.WAITING_BINDING) {
const flow = normalizeKuaishouCloudFlow(
parseTaskContext(workingTask).kuaishouCloudFulfillment,
)
if (!isKuaishouCloudMockTask(workingTask)) {
assertBoundUidMatchesExpected(workingTask, flow, {
errorCodePrefix: 'claim_kuaishou_cloud_redeem',
})
} else {
assertClaimExpectedUidReady(workingTask)
}
const confirmed = await updateTask(workingTask.id, {
task_status: TASK_STATUS.ROLE_CONFIRMED,
role_id: flow.binding.roleId || workingTask.role_id || '',
role_name: flow.binding.roleName || workingTask.role_name || '',
user_action_status: TASK_STATUS.ROLE_CONFIRMED,
role_confirmed_at: now,
last_error: '',
updated_at: now,
})
if (confirmed) {
workingTask = confirmed
}
await createTaskEvent(
workingTask.id,
'kuaishou_cloud_role_confirmed',
{
expectedUid: getClaimIdentityFromTask(workingTask).expectedUid,
vnPhone: flow.binding.vnPhone,
roleName: flow.binding.roleName,
roleId: flow.binding.roleId,
source: 'claim_page_one_click_redeem',
},
now,
)
}
const flowBeforeRedeem = normalizeKuaishouCloudFlow(
parseTaskContext(workingTask).kuaishouCloudFulfillment,
)
if (!isKuaishouCloudMockTask(workingTask)) {
assertBoundUidMatchesExpected(workingTask, flowBeforeRedeem, {
errorCodePrefix: 'claim_kuaishou_cloud_redeem',
})
} else {
assertClaimExpectedUidReady(workingTask)
}
const lockedTask = await updateTaskStatusIfCurrent(workingTask.id, TASK_STATUS.ROLE_CONFIRMED, {
task_status: TASK_STATUS.REDEEMING,
user_action_status: 'not_required',
last_error: '',
updated_at: now,
})
if (!lockedTask) {
return getKuaishouCloudClaimDetail(token)
}
if (isKuaishouCloudMockTask(lockedTask)) {
await completeMockKuaishouCloudClaimTask(lockedTask, now)
return getKuaishouCloudClaimDetail(token)
}
try {
await dispatchKuaishouCloudFulfillmentTask(lockedTask, {
source: 'claim_page_redeem',
actor: { source: 'claim_page' },
autoFinalize: true,
})
} catch (error) {
const latestTask = await getTaskById(lockedTask.id)
const latestStatus = latestTask ? normalizeTaskStatus(latestTask.task_status) : ''
if (latestTask && latestStatus !== TASK_STATUS.REDEEMING) {
if (
latestStatus === TASK_STATUS.DISPATCHED_PENDING_RETURN ||
latestStatus === TASK_STATUS.COMPLETED ||
latestStatus === TASK_STATUS.REDEEMED
) {
return getKuaishouCloudClaimDetail(token)
}
throw error
}
const message = error instanceof Error ? error.message : '兑换请求提交失败,请联系客服处理'
await updateTask(lockedTask.id, {
task_status: TASK_STATUS.MANUAL_REVIEW,
user_action_status: 'not_required',
last_error: message,
result_code: 'kuaishou_cloud_redeem_failed',
result_message: message,
updated_at: nowIso(),
})
await createTaskEvent(
lockedTask.id,
'kuaishou_cloud_redeem_failed',
{
source: 'claim_page_redeem',
errorMessage: message,
},
nowIso(),
)
throw error
}
return getKuaishouCloudClaimDetail(token)
}
function isKuaishouCloudMockTask(task: Partial<TaskRow> | null | undefined) {
return isKuaishouCloudMockContext(parseTaskContext(task))
}
function isKuaishouCloudMockContext(context: JsonObject = {}) {
const flow =
asJsonObject(context.kuaishouCloudFulfillment)
const mock = flow.mock && typeof flow.mock === 'object' ? flow.mock : context.mock
return Boolean(mock && typeof mock === 'object' && mock.enabled === true)
}
function buildMockVerifiedKuaishouCloudFlow(value: unknown, timestamp: string) {
const flow = normalizeKuaishouCloudFlow(value)
const source = flow as JsonObject
const bindUrl =
flow.binding.bindUrl ||
`https://example.com/mock-kuaishou-cloud-bind?task=mock&ts=${encodeURIComponent(timestamp)}`
const vnPhone = flow.binding.vnPhone || '13800000000'
const roleName = flow.binding.roleName || flow.role.name || '测试角色'
const roleId = flow.binding.roleId || flow.role.rid || '10001'
return {
...flow,
mock: {
...(asJsonObject(source.mock)),
enabled: true,
},
binding: {
...flow.binding,
prepareStatus: 'ready',
vnId: flow.binding.vnId || 900001,
vnPhone,
bindUrl,
bindPreparedAt: flow.binding.bindPreparedAt || timestamp,
bindExpiresAt: flow.binding.bindExpiresAt || addHours(timestamp, 24),
roleName,
roleId,
},
role: {
...flow.role,
status: 'ready',
name: roleName,
rid: roleId,
refreshedAt: flow.role.refreshedAt || timestamp,
errorMessage: '',
rawInfo: flow.role.rawInfo || {
mock: true,
},
},
}
}
async function completeMockKuaishouCloudClaimTask(task: TaskRow, timestamp: string) {
const taskContext = parseTaskContext(task)
const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment)
const source = flow as JsonObject
const nextFlow = {
...flow,
mock: {
...(asJsonObject(source.mock)),
enabled: true,
},
dispatch: {
...flow.dispatch,
status: 'success',
dispatchAt: timestamp,
dispatchBy: {
source: 'claim_page_mock',
},
note: '开发 mock 已模拟发货成功',
items: flow.deliveryItems,
},
returnNumber: {
...flow.returnNumber,
status: 'success',
returnedAt: timestamp,
returnedBy: {
source: 'claim_page_mock',
},
},
consume: {
...flow.consume,
status: 'success',
consumedAt: timestamp,
errorMessage: '',
},
}
await updateTask(task.id, {
task_status: TASK_STATUS.COMPLETED,
delivery_status: 'success',
result_code: 'mock_success',
result_message: '开发 mock 已模拟兑换成功',
user_action_status: 'not_required',
last_error: '',
context_json: JSON.stringify({
...taskContext,
kuaishouCloudFulfillment: nextFlow,
}),
redeemed_at: timestamp,
updated_at: timestamp,
})
await createTaskEvent(
task.id,
'kuaishou_cloud_mock_redeemed',
{
source: 'claim_page_mock',
deliveryItems: flow.deliveryItems,
},
timestamp,
)
}
@@ -1,8 +1,8 @@
import {
normalizeKuaishouCloudFlow,
prepareKuaishouCloudFulfillmentTask,
refreshKuaishouCloudTaskRoleInfo,
} from '../fulfillment/kuaishou-cloud/index.js'
prepareFulfillmentBinding,
refreshFulfillmentRole,
} from '../fulfillment/executors/registry.js'
import { normalizeKuaishouCloudFlow } from '../fulfillment/kuaishou-cloud/index.js'
import type { TaskRow } from '../../types/repository/rows.js'
import { parseTaskContext as parseTaskContextValue } from '../../utils/task-json.js'
@@ -12,11 +12,11 @@ export async function syncKuaishouCloudRoleInfo(task: TaskRow) {
if (!flow.binding.vnId || !flow.binding.vnKey || !flow.binding.vnPhone) {
if (flow.ticket.status === 'verified' && flow.binding.prepareStatus !== 'ready') {
try {
const prepared = await prepareKuaishouCloudFulfillmentTask(task, {
const prepared = await prepareFulfillmentBinding(task, {
source: 'claim_page_retry_binding_prepare',
actor: { source: 'system' },
})
return prepared.task
return prepared?.task || task
} catch {
return task
}
@@ -28,13 +28,13 @@ export async function syncKuaishouCloudRoleInfo(task: TaskRow) {
try {
// 只走 refresh:内部会按间隔 probe 绑链,必要时再查 bind_info
// 不再 forceProbe + 前置 probe,避免同轮重复打上游
const result = await refreshKuaishouCloudTaskRoleInfo(task, {
const result = await refreshFulfillmentRole(task, {
source: 'claim_page_polling',
actor: { source: 'system' },
recordEvent: false,
forceProbe: false,
})
return result.task || task
return result?.task || task
} catch {
return task
}
@@ -3,9 +3,20 @@ import {
shouldEnsureKuaishouCloudClaimLink,
TASK_STATUS,
} from '../../../domain/task-status.js'
import { ensureTaskClaimLink } from '../kuaishou-cloud/index.js'
import {
confirmKuaishouCloudTaskRole,
dispatchKuaishouCloudFulfillmentTask,
ensureTaskClaimLink,
prepareKuaishouCloudFulfillmentTask,
rebindKuaishouCloudTaskRole,
redeemKuaishouCloudTask,
refreshKuaishouCloudTaskRoleInfo,
returnKuaishouCloudFulfillmentTask,
} from '../kuaishou-cloud/index.js'
import {
FULFILLMENT_EXECUTOR_KEYS,
type FulfillmentActionOptions,
type FulfillmentActionResult,
type FulfillmentDeliveryLink,
type FulfillmentExecutor,
type FulfillmentPrepareDeps,
@@ -16,6 +27,13 @@ export const kuaishouCloudExecutor: FulfillmentExecutor = {
key: FULFILLMENT_EXECUTOR_KEYS.KUAISHOU_CLOUD,
preparePaidTask,
resolveDeliveryLink,
prepareBinding,
rebindRole,
refreshRole,
confirmRole,
redeemTask,
dispatchTask,
returnNumber,
}
async function preparePaidTask(
@@ -71,3 +89,115 @@ async function resolveDeliveryLink(task: TaskRow): Promise<FulfillmentDeliveryLi
return null
}
}
async function prepareBinding(
task: TaskRow,
options: FulfillmentActionOptions = {},
): Promise<FulfillmentActionResult> {
const result = await prepareKuaishouCloudFulfillmentTask(task, {
source: options.source || 'executor_prepare_binding',
actor: options.actor,
force: options.force === true,
})
return {
task: result.task || task,
claimUrl: result.claimUrl,
token: result.token,
flow: result.flow,
}
}
async function rebindRole(
task: TaskRow,
options: FulfillmentActionOptions = {},
): Promise<FulfillmentActionResult> {
const result = await rebindKuaishouCloudTaskRole(task, {
source: options.source || 'executor_rebind_role',
actor: options.actor,
})
return {
task: result.task || task,
claimUrl: result.claimUrl,
token: result.token,
}
}
async function refreshRole(
task: TaskRow,
options: FulfillmentActionOptions = {},
): Promise<FulfillmentActionResult> {
const result = await refreshKuaishouCloudTaskRoleInfo(task, {
source: options.source || 'executor_refresh_role',
actor: options.actor,
recordEvent: options.recordEvent !== false,
forceProbe: options.forceProbe === true,
})
return {
task: result.task || task,
roleInfo: result.roleInfo,
}
}
async function confirmRole(
task: TaskRow,
options: FulfillmentActionOptions = {},
): Promise<FulfillmentActionResult> {
const result = await confirmKuaishouCloudTaskRole(task, {
source: options.source || 'executor_confirm_role',
actor: options.actor,
errorCodePrefix: options.errorCodePrefix,
forceProbe: options.forceProbe !== false,
})
return {
task: result.task,
alreadySettled: result.alreadySettled,
}
}
async function redeemTask(
task: TaskRow,
options: FulfillmentActionOptions = {},
): Promise<FulfillmentActionResult> {
const result = await redeemKuaishouCloudTask(task, {
source: options.source || 'executor_redeem',
actor: options.actor,
autoFinalize: options.autoFinalize !== false,
errorCodePrefix: options.errorCodePrefix,
})
return {
task: result.task,
alreadySettled: result.alreadySettled,
concurrentSkipped: result.concurrentSkipped,
}
}
async function dispatchTask(
task: TaskRow,
options: FulfillmentActionOptions = {},
): Promise<FulfillmentActionResult> {
const result = await dispatchKuaishouCloudFulfillmentTask(task, {
source: options.source || 'executor_dispatch',
actor: options.actor,
autoFinalize: options.autoFinalize === true,
})
return {
task: result.task || task,
flow: result.flow,
}
}
async function returnNumber(
task: TaskRow,
options: FulfillmentActionOptions = {},
): Promise<FulfillmentActionResult> {
const result = await returnKuaishouCloudFulfillmentTask(task, {
source: options.source || 'executor_return_number',
actor: options.actor,
// 未显式传入时保持默认 true(发货收尾);admin 清理会传 false
consumeIndustryVoucher: options.consumeIndustryVoucher,
})
return {
task: result.task || task,
flow: result.flow,
}
}
@@ -0,0 +1,70 @@
import assert from 'node:assert/strict'
import test from 'node:test'
import {
getFulfillmentExecutor,
prepareFulfillmentBinding,
redeemFulfillmentTask,
dispatchFulfillmentTask,
confirmFulfillmentRole,
rebindFulfillmentRole,
refreshFulfillmentRole,
returnFulfillmentNumber,
} from './registry.js'
import { FULFILLMENT_EXECUTOR_KEYS } from './types.js'
import type { TaskRow } from '../../../types/repository/rows.js'
function makeTask(executorKey: string): TaskRow {
return {
id: 1,
executor_key: executorKey,
} as TaskRow
}
test('getFulfillmentExecutor 映射 lewan / industry / feifei / manual', () => {
assert.equal(
getFulfillmentExecutor(FULFILLMENT_EXECUTOR_KEYS.KUAISHOU_CLOUD)?.key,
FULFILLMENT_EXECUTOR_KEYS.KUAISHOU_CLOUD,
)
assert.equal(
getFulfillmentExecutor(FULFILLMENT_EXECUTOR_KEYS.KUAISHOU_INDUSTRY)?.key,
FULFILLMENT_EXECUTOR_KEYS.KUAISHOU_CLOUD,
)
assert.equal(
getFulfillmentExecutor(FULFILLMENT_EXECUTOR_KEYS.KUAISHOU_FEIFEI)?.key,
FULFILLMENT_EXECUTOR_KEYS.KUAISHOU_FEIFEI,
)
assert.equal(
getFulfillmentExecutor(FULFILLMENT_EXECUTOR_KEYS.MANUAL_DISPATCH)?.key,
FULFILLMENT_EXECUTOR_KEYS.MANUAL_DISPATCH,
)
assert.equal(getFulfillmentExecutor('unknown'), null)
})
test('lewan executor 暴露完整履约动作', () => {
const executor = getFulfillmentExecutor(FULFILLMENT_EXECUTOR_KEYS.KUAISHOU_CLOUD)
assert.ok(executor)
assert.equal(typeof executor.prepareBinding, 'function')
assert.equal(typeof executor.rebindRole, 'function')
assert.equal(typeof executor.refreshRole, 'function')
assert.equal(typeof executor.confirmRole, 'function')
assert.equal(typeof executor.redeemTask, 'function')
assert.equal(typeof executor.dispatchTask, 'function')
assert.equal(typeof executor.returnNumber, 'function')
})
test('feifei / manual 不支持 lewan 专属动作时 registry 返回 null', async () => {
const feifei = makeTask(FULFILLMENT_EXECUTOR_KEYS.KUAISHOU_FEIFEI)
const manual = makeTask(FULFILLMENT_EXECUTOR_KEYS.MANUAL_DISPATCH)
assert.equal(await prepareFulfillmentBinding(feifei), null)
assert.equal(await rebindFulfillmentRole(feifei), null)
assert.equal(await refreshFulfillmentRole(feifei), null)
assert.equal(await confirmFulfillmentRole(feifei), null)
assert.equal(await redeemFulfillmentTask(feifei), null)
assert.equal(await dispatchFulfillmentTask(feifei), null)
assert.equal(await returnFulfillmentNumber(feifei), null)
assert.equal(await redeemFulfillmentTask(manual), null)
assert.equal(await dispatchFulfillmentTask(manual), null)
})
@@ -6,6 +6,7 @@ import {
FULFILLMENT_EXECUTOR_KEYS,
isManualDispatchExecutor,
normalizeExecutorKey,
type FulfillmentActionOptions,
type FulfillmentDeliveryLink,
type FulfillmentExecutor,
type FulfillmentPrepareDeps,
@@ -53,3 +54,73 @@ export async function resolveFulfillmentDeliveryLink(
return executor.resolveDeliveryLink(task)
}
async function runExecutorAction(
task: TaskRow,
method:
| 'prepareBinding'
| 'rebindRole'
| 'refreshRole'
| 'confirmRole'
| 'redeemTask'
| 'dispatchTask'
| 'returnNumber',
options: FulfillmentActionOptions = {},
) {
const executor = getFulfillmentExecutor(task.executor_key)
const handler = executor?.[method]
if (!handler) {
return null
}
return handler(task, options)
}
export function prepareFulfillmentBinding(
task: TaskRow,
options: FulfillmentActionOptions = {},
) {
return runExecutorAction(task, 'prepareBinding', options)
}
export function rebindFulfillmentRole(
task: TaskRow,
options: FulfillmentActionOptions = {},
) {
return runExecutorAction(task, 'rebindRole', options)
}
export function refreshFulfillmentRole(
task: TaskRow,
options: FulfillmentActionOptions = {},
) {
return runExecutorAction(task, 'refreshRole', options)
}
export function confirmFulfillmentRole(
task: TaskRow,
options: FulfillmentActionOptions = {},
) {
return runExecutorAction(task, 'confirmRole', options)
}
export function redeemFulfillmentTask(
task: TaskRow,
options: FulfillmentActionOptions = {},
) {
return runExecutorAction(task, 'redeemTask', options)
}
export function dispatchFulfillmentTask(
task: TaskRow,
options: FulfillmentActionOptions = {},
) {
return runExecutorAction(task, 'dispatchTask', options)
}
export function returnFulfillmentNumber(
task: TaskRow,
options: FulfillmentActionOptions = {},
) {
return runExecutorAction(task, 'returnNumber', options)
}
@@ -31,6 +31,25 @@ export type FulfillmentPrepareDeps = {
nowIso: () => string
}
export type FulfillmentActionOptions = {
source?: string
actor?: unknown
autoFinalize?: boolean
errorCodePrefix?: string
/**
* 退号时是否顺带核销行业电子凭证。
* - true/默认:发货收尾(claim autoFinalize
* - false:仅清理虚拟号(admin 退号),与核销无关
*/
consumeIndustryVoucher?: boolean
[key: string]: unknown
}
export type FulfillmentActionResult = {
task: TaskRow
[key: string]: unknown
}
export type FulfillmentExecutor = {
key: FulfillmentExecutorKey
preparePaidTask?: (
@@ -38,6 +57,44 @@ export type FulfillmentExecutor = {
deps: FulfillmentPrepareDeps,
) => Promise<TaskRow | null>
resolveDeliveryLink?: (task: TaskRow) => Promise<FulfillmentDeliveryLink | null>
/** lewan:准备绑定资源(虚拟号 / bindUrl) */
prepareBinding?: (
task: TaskRow,
options?: FulfillmentActionOptions,
) => Promise<FulfillmentActionResult>
/** lewan:换绑角色 */
rebindRole?: (
task: TaskRow,
options?: FulfillmentActionOptions,
) => Promise<FulfillmentActionResult>
/** lewan:刷新角色信息 */
refreshRole?: (
task: TaskRow,
options?: FulfillmentActionOptions,
) => Promise<FulfillmentActionResult>
/** lewan:确认角色(UID 匹配后落 ROLE_CONFIRMED */
confirmRole?: (
task: TaskRow,
options?: FulfillmentActionOptions,
) => Promise<FulfillmentActionResult>
/** lewan:一键兑换(状态机 + dispatch */
redeemTask?: (
task: TaskRow,
options?: FulfillmentActionOptions,
) => Promise<FulfillmentActionResult>
/** lewan:后台/系统发货(不经 claim 一键兑换) */
dispatchTask?: (
task: TaskRow,
options?: FulfillmentActionOptions,
) => Promise<FulfillmentActionResult>
/**
* lewan:退还虚拟号。
* 默认可顺带核销;admin 清理请传 consumeIndustryVoucher: false。
*/
returnNumber?: (
task: TaskRow,
options?: FulfillmentActionOptions,
) => Promise<FulfillmentActionResult>
}
export function normalizeExecutorKey(value: unknown): FulfillmentExecutorKey {
@@ -0,0 +1,106 @@
import { createTaskEvent } from '../../../repositories/task-event-repo.js'
import { updateTask } from '../../../repositories/task-repo.js'
import {
TASK_STATUS,
isKuaishouCloudRoleConfirmSettledStatus,
} from '../../../domain/task-status.js'
import { createHttpError } from '../../../utils/http.js'
import { nowIso } from '../../../utils/time.js'
import type { TaskRow } from '../../../types/repository/rows.js'
import {
assertBoundUidMatchesExpected,
assertClaimExpectedUidReady,
} from '../../claim/claim-identity.js'
import { isKuaishouCloudTask, normalizeKuaishouCloudFlow, type JsonObject } from './domain.js'
import { isKuaishouCloudMockContext, isKuaishouCloudMockTask } from './mock-helpers.js'
import { refreshKuaishouCloudTaskRoleInfo } from './refresh-role-info.js'
import { normalizeActor, parseTaskContext } from './task-context.js'
export type ConfirmKuaishouCloudRoleResult = {
task: TaskRow
alreadySettled: boolean
}
/**
* lewan 确认角色用例:刷新角色 → UID 闸门 → ROLE_CONFIRMED。
* 一键兑换路径里 WAITING_BINDING 的自动晋升在 redeem 内完成,不必先调本函数。
*/
export async function confirmKuaishouCloudTaskRole(
task: TaskRow,
options: JsonObject = {},
): Promise<ConfirmKuaishouCloudRoleResult> {
if (!isKuaishouCloudTask(task)) {
throw createHttpError('当前任务不是 kuaishou-lewan 履约任务', {
statusCode: 409,
errorCode: 'kuaishou_cloud_task_invalid',
})
}
const now = String(options.now || '').trim() || nowIso()
const source = String(options.source || 'system_confirm_role').trim() || 'system_confirm_role'
const actor = normalizeActor(options.actor)
const errorCodePrefix =
String(options.errorCodePrefix || 'kuaishou_cloud_confirm').trim() || 'kuaishou_cloud_confirm'
if (isKuaishouCloudRoleConfirmSettledStatus(task.task_status)) {
return { task, alreadySettled: true }
}
const expectedUid = assertClaimExpectedUidReady(task)
const mockMode = isKuaishouCloudMockTask(task) || isKuaishouCloudMockContext(parseTaskContext(task))
const refreshed = mockMode
? { task }
: await refreshKuaishouCloudTaskRoleInfo(task, {
source,
actor,
recordEvent: options.recordEvent === true,
forceProbe: options.forceProbe !== false,
})
const flow = normalizeKuaishouCloudFlow(
parseTaskContext(refreshed.task).kuaishouCloudFulfillment,
)
if (!flow.binding.vnPhone || !flow.binding.roleName || !flow.binding.roleId) {
throw createHttpError('角色信息还未刷新到系统,请完成绑定后稍等片刻再试', {
statusCode: 409,
errorCode: `${errorCodePrefix}_role_not_ready`,
})
}
if (!mockMode) {
assertBoundUidMatchesExpected(refreshed.task, flow, {
errorCodePrefix,
})
}
const updatedTask =
(await updateTask(task.id, {
task_status: TASK_STATUS.ROLE_CONFIRMED,
role_id: flow.binding.roleId,
role_name: flow.binding.roleName,
user_action_status: TASK_STATUS.ROLE_CONFIRMED,
role_confirmed_at: now,
last_error: '',
updated_at: now,
})) ||
refreshed.task ||
task
await createTaskEvent(
task.id,
'kuaishou_cloud_role_confirmed',
{
expectedUid,
vnPhone: flow.binding.vnPhone,
roleName: flow.binding.roleName,
roleId: flow.binding.roleId,
source,
actor,
},
now,
)
return { task: updatedTask, alreadySettled: false }
}
@@ -0,0 +1,131 @@
import { createTaskEvent } from '../../../repositories/task-event-repo.js'
import { updateTask } from '../../../repositories/task-repo.js'
import { createHttpError } from '../../../utils/http.js'
import { nowIso } from '../../../utils/time.js'
import { TASK_STATUS } from '../../../domain/task-status.js'
import type { TaskRow } from '../../../types/repository/rows.js'
import type { JsonObject } from './domain.js'
import type { DispatchStockResult } from './dispatch-types.js'
export async function markKuaishouCloudDispatchFailed(
task: TaskRow,
taskContext: JsonObject,
flow: JsonObject,
error: unknown,
options: JsonObject = {},
) {
const now = nowIso()
const errorMessage = resolveErrorMessage(error)
const errorCode = resolveErrorCode(error) || 'kuaishou_cloud_dispatch_failed'
const failureContext = resolveDispatchFailureContext(error)
const stockResult = isPlainObject(failureContext.stockResult)
? (failureContext.stockResult as Partial<DispatchStockResult>)
: null
const nextContext = {
...taskContext,
kuaishouCloudFulfillment: {
...flow,
dispatch: {
...(isPlainObject(flow.dispatch) ? flow.dispatch : {}),
status: 'failed',
failedAt: now,
failedStage: String(failureContext.stage || '').trim(),
errorCode,
errorMessage,
note: buildDispatchFailedMessage(errorMessage),
},
purchase: {
...(isPlainObject(flow.purchase) ? flow.purchase : {}),
...(stockResult
? {
usedKnapsack: stockResult.usedKnapsack,
purchaseTriggered: stockResult.purchaseTriggered,
assetBefore: stockResult.assetBefore,
assetAfter: stockResult.assetAfter,
items: stockResult.items,
}
: {}),
},
},
}
await updateTask(task.id, {
task_status: TASK_STATUS.MANUAL_REVIEW,
user_action_status: 'not_required',
last_error: errorMessage,
result_code: errorCode,
result_message: buildDispatchFailedMessage(errorMessage),
context_json: JSON.stringify(nextContext),
updated_at: now,
})
await createTaskEvent(
task.id,
'kuaishou_cloud_dispatch_failed',
{
source: String(options.source || 'system').trim() || 'system',
actor: options.actor,
errorCode,
errorMessage,
failureContext,
},
now,
)
}
export function enrichCloudtentaclesDispatchError(error: unknown, context: JsonObject) {
const currentError =
error instanceof Error
? (error as Error & { context?: unknown })
: (createHttpError(resolveErrorMessage(error), {
statusCode: 500,
errorCode: 'kuaishou_cloud_dispatch_failed',
}) as Error & { context?: unknown })
const currentContext = isPlainObject(currentError.context) ? currentError.context : {}
currentError.context = {
...currentContext,
dispatchFailure: {
...(isPlainObject((currentContext as JsonObject).dispatchFailure)
? ((currentContext as JsonObject).dispatchFailure as JsonObject)
: {}),
...context,
},
}
return currentError
}
export function resolveDispatchFailureContext(error: unknown): JsonObject {
if (!error || typeof error !== 'object') {
return {}
}
const context = (error as { context?: unknown }).context
if (!isPlainObject(context)) {
return {}
}
const dispatchFailure = context.dispatchFailure
return isPlainObject(dispatchFailure) ? dispatchFailure : {}
}
export function resolveErrorCode(error: unknown) {
if (!error || typeof error !== 'object') {
return ''
}
return String((error as { errorCode?: unknown }).errorCode || '').trim()
}
export function resolveErrorMessage(error: unknown) {
return error instanceof Error ? error.message : String(error || 'CloudTentacles 发货失败')
}
export function buildDispatchFailedMessage(message: unknown) {
const reason = String(message || '').trim() || '未知错误'
return `CloudTentacles 发货失败:${reason}`
}
function isPlainObject(value: unknown): value is JsonObject {
return Boolean(value) && typeof value === 'object' && !Array.isArray(value)
}
@@ -0,0 +1,247 @@
import { createTaskEvent } from '../../../repositories/task-event-repo.js'
import { updateTask } from '../../../repositories/task-repo.js'
import { createHttpError } from '../../../utils/http.js'
import { nowIso } from '../../../utils/time.js'
import { TASK_STATUS, normalizeTaskStatus } from '../../../domain/task-status.js'
import type { TaskRow } from '../../../types/repository/rows.js'
import {
isIndustryEVoucherTask,
isKuaishouCloudTask,
maskCode,
maskPhone,
normalizeKuaishouCloudFlow,
type JsonObject,
} from './domain.js'
import { resolvePersistedCloudtentaclesContextBySourceKeys } from './cloudtentacles-context.js'
import { syncKuaishouCloudRoleInfoBeforeDispatch } from './dispatch-role-sync.js'
import { markKuaishouCloudDispatchFailed } from './dispatch-failure.js'
import {
prepareStockAndDispatch,
resolveCloudSkuDispatchLockKeys,
withCloudSkuDispatchLocks,
} from './dispatch-stock.js'
import {
resolveDispatchDeliveryItems,
type DispatchResultItem,
type DispatchStockResult,
} from './dispatch-types.js'
import { returnKuaishouCloudFulfillmentTask } from './return-fulfillment.js'
import { normalizeActor, parseTaskContext } from './task-context.js'
/**
* lewan 自动/后台发货主用例:
* UID 闸门同步角色 → 库存与采购 → cloudtentacles 下发 → 可选自动退号收尾。
*/
export async function dispatchKuaishouCloudFulfillmentTask(
task: TaskRow,
options: JsonObject = {},
) {
if (!isKuaishouCloudTask(task)) {
throw createHttpError('当前任务不是 kuaishou-lewan 履约任务', {
statusCode: 409,
errorCode: 'kuaishou_cloud_task_invalid',
})
}
const actor = normalizeActor(options.actor)
const now = nowIso()
const taskContext = parseTaskContext(task)
const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment)
const cloudContext = resolvePersistedCloudtentaclesContextBySourceKeys([
flow.binding.resolvedSourceKey,
...flow.binding.cloudSourceKeys,
])
if (!flow.binding.skuId || !flow.binding.vnId || !flow.binding.vnPhone) {
throw createHttpError('当前任务还没有准备好绑定资源,请先完成绑定资源准备', {
statusCode: 409,
errorCode: 'kuaishou_cloud_not_prepared',
})
}
const persistedTicketCode = String(flow.ticket.code || '').trim()
const voucherContext = isPlainObject(taskContext.kuaishouIndustryVoucher)
? taskContext.kuaishouIndustryVoucher
: {}
const industryVoucherCode = String(
voucherContext.voucherCode || voucherContext.eticketId || '',
).trim()
const hasIndustryVoucherForDispatch =
Boolean(industryVoucherCode) || isIndustryEVoucherTask(task)
const resolvedTicketCode = persistedTicketCode || industryVoucherCode
if (!resolvedTicketCode && !hasIndustryVoucherForDispatch) {
throw createHttpError('旧快手小店核销流程已停用,请改用行业电子凭证处理', {
statusCode: 409,
errorCode: 'kuaishou_cloud_missing_ticket_code',
})
}
if (
flow.dispatch.status === 'success' &&
normalizeTaskStatus(task.task_status) === TASK_STATUS.DISPATCHED_PENDING_RETURN
) {
return { task, flow }
}
const synced = await syncKuaishouCloudRoleInfoBeforeDispatch(task, {
now,
actor,
source: options.source || 'system_before_dispatch',
cloudContext,
taskContext,
flow,
})
const syncedFlow = synced.flow
const syncedTaskContext = synced.taskContext
const deliveryItems = resolveDispatchDeliveryItems(syncedFlow as JsonObject)
if (deliveryItems.length === 0) {
throw createHttpError('当前任务缺少 cloud 发货物品配置', {
statusCode: 409,
errorCode: 'kuaishou_cloud_missing_delivery_items',
})
}
let dispatchResults: DispatchResultItem[]
let stockResult: DispatchStockResult
try {
;({ dispatchResults, stockResult } = await withCloudSkuDispatchLocks(
resolveCloudSkuDispatchLockKeys(cloudContext.resolvedSourceKey, deliveryItems),
() =>
prepareStockAndDispatch({
task,
flow: syncedFlow as JsonObject,
cloudContext,
deliveryItems,
}),
))
} catch (error) {
await markKuaishouCloudDispatchFailed(
task,
syncedTaskContext as JsonObject,
syncedFlow as JsonObject,
error,
{
actor,
source: options.source || 'system',
},
)
throw error
}
const firstDispatchResult: DispatchResultItem = dispatchResults[0] || {
cloudSkuId: syncedFlow.binding.skuId,
cloudSkuName: syncedFlow.binding.skuName,
quantity: 1,
unitIndex: 1,
sendType: 0,
note: '',
responseMessage: '',
}
const dispatchSummary =
dispatchResults.length > 1
? `cloudtentacles 发货成功,共 ${dispatchResults.length}`
: String(
firstDispatchResult.responseMessage ||
firstDispatchResult.note ||
'cloudtentacles 发货成功',
).trim()
const nextContext = {
...syncedTaskContext,
kuaishouCloudFulfillment: {
...syncedFlow,
ticket: {
...syncedFlow.ticket,
code: resolvedTicketCode,
capturedAt: resolvedTicketCode
? syncedFlow.ticket.capturedAt || now
: syncedFlow.ticket.capturedAt,
capturedBy: syncedFlow.ticket.capturedBy,
},
dispatch: {
...syncedFlow.dispatch,
status: 'success',
dispatchAt: now,
dispatchBy: actor,
sendType: Number(firstDispatchResult.sendType || 0) || 0,
note: dispatchSummary,
items: dispatchResults,
},
purchase: {
...syncedFlow.purchase,
usedKnapsack: stockResult.usedKnapsack,
purchaseTriggered: stockResult.purchaseTriggered,
assetBefore: stockResult.assetBefore,
assetAfter: stockResult.assetAfter,
purchaseAt: stockResult.purchaseTriggered
? now
: syncedFlow.purchase.purchaseAt,
items: stockResult.items,
},
},
}
let updatedTask = await updateTask(task.id, {
task_status: TASK_STATUS.DISPATCHED_PENDING_RETURN,
delivery_status: 'delivered',
result_code: 'kuaishou_cloud_dispatched',
result_message: dispatchSummary,
user_action_status: 'not_required',
last_error: '',
context_json: JSON.stringify(nextContext),
updated_at: now,
})
if (!updatedTask) {
throw createHttpError('kuaishou-lewan 发货状态更新失败', {
statusCode: 500,
errorCode: 'kuaishou_cloud_dispatch_update_failed',
})
}
await createTaskEvent(
task.id,
'kuaishou_cloud_dispatched',
{
source: String(options.source || 'system').trim() || 'system',
ticketCodeMasked: maskCode(resolvedTicketCode),
skuId: syncedFlow.binding.skuId,
vnId: syncedFlow.binding.vnId,
vnPhoneMasked: maskPhone(syncedFlow.binding.vnPhone),
sendType: firstDispatchResult.sendType,
note: dispatchSummary,
deliveryItems,
dispatchResults,
stockResult,
actor,
},
now,
)
const shouldAutoFinalize =
options.autoFinalize === true &&
normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment).returnNumber
.autoReturnEnabled === true
if (shouldAutoFinalize) {
const finalizeResult = await returnKuaishouCloudFulfillmentTask(updatedTask, {
actor,
source: options.source || 'system_auto_finalize',
// 发货后自动收尾:退号并尝试核销(与 admin 清理退号不同)
consumeIndustryVoucher: true,
})
updatedTask = finalizeResult.task
}
return {
task: updatedTask,
flow: normalizeKuaishouCloudFlow(
parseTaskContext(updatedTask).kuaishouCloudFulfillment,
),
}
}
function isPlainObject(value: unknown): value is JsonObject {
return Boolean(value) && typeof value === 'object' && !Array.isArray(value)
}
@@ -0,0 +1,366 @@
import { createTaskEvent } from '../../../repositories/task-event-repo.js'
import { createHttpError } from '../../../utils/http.js'
import { nowIso } from '../../../utils/time.js'
import { notifyKuaishouCloudAssetNotEnough } from '../../notification/domain-notifications.js'
import {
buyCloudtentaclesSku,
getCloudtentaclesAsset,
listCloudtentaclesSku,
useCloudtentaclesSku,
} from '../../platforms/cloudtentacles/catalog-service.js'
import { getCloudtentaclesKnapsack } from '../../platforms/cloudtentacles/knapsack-service.js'
import type { TaskRow } from '../../../types/repository/rows.js'
import { maskPhone, type JsonObject } from './domain.js'
import {
enrichCloudtentaclesDispatchError,
resolveErrorCode,
resolveErrorMessage,
} from './dispatch-failure.js'
import type {
DispatchDeliveryItem,
DispatchResultItem,
DispatchStockItem,
DispatchStockResult,
} from './dispatch-types.js'
const cloudSkuDispatchLocks = new Map<string, Promise<void>>()
/**
* 计算库存缺口 → 必要时采购 → 按数量逐次 use SKU 下发。
* 调用方应在 withCloudSkuDispatchLocks 内执行,避免同账号同 SKU 并发抢库存。
*/
export async function prepareStockAndDispatch({
task,
flow,
cloudContext,
deliveryItems,
}: {
task: TaskRow
flow: JsonObject
cloudContext: JsonObject
deliveryItems: DispatchDeliveryItem[]
}): Promise<{ dispatchResults: DispatchResultItem[]; stockResult: DispatchStockResult }> {
const [knapsack, skuList] = await Promise.all([
getCloudtentaclesKnapsack(cloudContext),
listCloudtentaclesSku(cloudContext),
])
const skuItems = Array.isArray(skuList.items) ? (skuList.items as JsonObject[]) : []
const knapsackItems = Array.isArray(knapsack.items) ? (knapsack.items as JsonObject[]) : []
const stockItems = buildDispatchStockItems(deliveryItems, {
skuItems,
knapsackItems,
})
const missingItems = stockItems.filter((item) => item.purchasedCount > 0)
let assetBefore = 0
let assetAfter = 0
let purchaseTriggered = false
const purchase = isPlainObject(flow.purchase) ? flow.purchase : {}
const binding = isPlainObject(flow.binding) ? flow.binding : {}
if (missingItems.length > 0) {
if (purchase.autoBuyEnabled === false) {
throw createHttpError('背包中没有现成库存,且当前配置未开启自动购买', {
statusCode: 409,
errorCode: 'kuaishou_cloud_auto_buy_disabled',
})
}
const missingSku = missingItems.find((item) => !findCloudSkuItem(skuItems, item.cloudSkuId))
if (missingSku) {
throw createHttpError(`cloudtentacles 未找到 SKU ${missingSku.cloudSkuId}`, {
statusCode: 404,
errorCode: 'kuaishou_cloud_sku_not_found',
})
}
const asset = await getCloudtentaclesAsset(cloudContext)
assetBefore = Number(asset.asset || 0) || 0
const targetPrice = missingItems.reduce((sum, item) => {
const skuItem = findCloudSkuItem(skuItems, item.cloudSkuId)
return sum + item.purchasedCount * Number(skuItem?.price || 0)
}, 0)
const requiredAsset = targetPrice + Number(purchase.minAssetReserve || 0)
if (assetBefore < requiredAsset) {
await notifyKuaishouCloudAssetNotEnough({
task,
flow,
assetBefore,
requiredAsset,
skuName: missingItems
.map((item) => item.cloudSkuName || `SKU ${item.cloudSkuId}`)
.join('、'),
})
throw createHttpError(`余额不足,当前 ${assetBefore},至少需要 ${requiredAsset}`, {
statusCode: 409,
errorCode: 'kuaishou_cloud_asset_not_enough',
})
}
for (const item of missingItems) {
await createTaskEvent(
task.id,
'cloudtentacles_sku_buy_started',
{
source: 'kuaishou_cloud_dispatch',
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
count: item.purchasedCount,
assetBefore,
},
nowIso(),
)
try {
const buyResult = await buyCloudtentaclesSku({
...cloudContext,
id: item.cloudSkuId,
count: item.purchasedCount,
})
await createTaskEvent(
task.id,
'cloudtentacles_sku_buy_succeeded',
{
source: 'kuaishou_cloud_dispatch',
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
count: item.purchasedCount,
responseMessage: buyResult.responseMessage,
},
nowIso(),
)
} catch (error) {
await createTaskEvent(
task.id,
'cloudtentacles_sku_buy_failed',
{
source: 'kuaishou_cloud_dispatch',
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
count: item.purchasedCount,
errorCode: resolveErrorCode(error),
errorMessage: resolveErrorMessage(error),
},
nowIso(),
)
throw enrichCloudtentaclesDispatchError(error, {
stage: 'purchase',
stockResult: buildCurrentStockResult({
stockItems,
purchaseTriggered,
assetBefore,
assetAfter,
}),
item,
})
}
}
purchaseTriggered = true
const assetResult = await getCloudtentaclesAsset(cloudContext)
assetAfter = Number(assetResult.asset || 0) || 0
}
const dispatchResults: DispatchResultItem[] = []
const vnId = Number(binding.vnId || 0) || 0
const vnPhone = String(binding.vnPhone || '').trim()
for (const item of deliveryItems) {
for (let index = 0; index < item.quantity; index += 1) {
await createTaskEvent(
task.id,
'cloudtentacles_sku_use_started',
{
source: 'kuaishou_cloud_dispatch',
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
unitIndex: index + 1,
quantity: item.quantity,
vnId,
vnPhoneMasked: maskPhone(vnPhone),
},
nowIso(),
)
let dispatchResult: Awaited<ReturnType<typeof useCloudtentaclesSku>>
try {
dispatchResult = await useCloudtentaclesSku({
...cloudContext,
id: item.cloudSkuId,
virtualNumberId: vnId,
phone: vnPhone,
})
} catch (error) {
await createTaskEvent(
task.id,
'cloudtentacles_sku_use_failed',
{
source: 'kuaishou_cloud_dispatch',
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
unitIndex: index + 1,
quantity: item.quantity,
vnId,
vnPhoneMasked: maskPhone(vnPhone),
errorCode: resolveErrorCode(error),
errorMessage: resolveErrorMessage(error),
},
nowIso(),
)
throw enrichCloudtentaclesDispatchError(error, {
stage: 'dispatch',
stockResult: buildCurrentStockResult({
stockItems,
purchaseTriggered,
assetBefore,
assetAfter,
}),
item,
unitIndex: index + 1,
vnId,
vnPhoneMasked: maskPhone(vnPhone),
})
}
await createTaskEvent(
task.id,
'cloudtentacles_sku_use_succeeded',
{
source: 'kuaishou_cloud_dispatch',
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
unitIndex: index + 1,
quantity: item.quantity,
vnId,
vnPhoneMasked: maskPhone(vnPhone),
sendType: Number(dispatchResult.sendType || 0) || 0,
note: String(dispatchResult.note || '').trim(),
responseMessage: String(dispatchResult.responseMessage || '').trim(),
},
nowIso(),
)
dispatchResults.push({
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
unitIndex: index + 1,
quantity: item.quantity,
sendType: Number(dispatchResult.sendType || 0) || 0,
note: String(dispatchResult.note || '').trim(),
responseMessage: String(dispatchResult.responseMessage || '').trim(),
})
}
}
return {
dispatchResults,
stockResult: {
usedKnapsack: stockItems.every((item) => item.purchasedCount <= 0),
purchaseTriggered,
assetBefore,
assetAfter,
items: stockItems,
},
}
}
export function buildDispatchStockItems(
deliveryItems: DispatchDeliveryItem[],
{
skuItems = [],
knapsackItems = [],
}: { skuItems?: JsonObject[]; knapsackItems?: JsonObject[] } = {},
): DispatchStockItem[] {
return deliveryItems.map((item) => {
const knapsackItem = findCloudSkuItem(knapsackItems, item.cloudSkuId)
const skuItem = findCloudSkuItem(skuItems, item.cloudSkuId)
const knapsackCount = Math.max(0, Number(knapsackItem?.count || 0) || 0)
return {
...item,
cloudSkuName: item.cloudSkuName || String(skuItem?.name || knapsackItem?.name || '').trim(),
requiredCount: item.quantity,
knapsackCount,
purchasedCount: Math.max(0, item.quantity - knapsackCount),
}
})
}
export function resolveCloudSkuDispatchLockKeys(
sourceKey: unknown,
deliveryItems: DispatchDeliveryItem[],
) {
const normalizedSourceKey = String(sourceKey || 'default').trim() || 'default'
return deliveryItems.map((item) => `${normalizedSourceKey}:${item.cloudSkuId}`)
}
export async function withCloudSkuDispatchLocks<T>(
keys: string[],
callback: () => Promise<T>,
): Promise<T> {
const normalizedKeys = [
...new Set(keys.map((key) => String(key || '').trim()).filter(Boolean)),
].sort()
async function run(index: number): Promise<T> {
const key = normalizedKeys[index]
if (!key) {
return callback()
}
return withCloudSkuDispatchLock(key, () => run(index + 1))
}
return run(0)
}
function buildCurrentStockResult({
stockItems,
purchaseTriggered,
assetBefore,
assetAfter,
}: {
stockItems: DispatchStockItem[]
purchaseTriggered: boolean
assetBefore: number
assetAfter: number
}): DispatchStockResult {
return {
usedKnapsack: stockItems.every((item) => item.purchasedCount <= 0),
purchaseTriggered,
assetBefore,
assetAfter,
items: stockItems,
}
}
function findCloudSkuItem(items: JsonObject[], cloudSkuId: number) {
return items.find((item) => Number(item.id || 0) === cloudSkuId) || null
}
async function withCloudSkuDispatchLock<T>(key: string, callback: () => Promise<T>): Promise<T> {
const previous = cloudSkuDispatchLocks.get(key) || Promise.resolve()
let release: () => void = () => {}
const current = new Promise<void>((resolve) => {
release = resolve
})
const queued = previous.catch(() => {}).then(() => current)
cloudSkuDispatchLocks.set(key, queued)
await previous.catch(() => {})
try {
return await callback()
} finally {
release()
if (cloudSkuDispatchLocks.get(key) === queued) {
cloudSkuDispatchLocks.delete(key)
}
}
}
function isPlainObject(value: unknown): value is JsonObject {
return Boolean(value) && typeof value === 'object' && !Array.isArray(value)
}
@@ -0,0 +1,69 @@
import type { JsonObject } from './domain.js'
/** lewan 发货单元:一次下发的 cloud SKU */
export type DispatchDeliveryItem = {
cloudSkuId: number
cloudSkuName: string
quantity: number
}
export type DispatchResultItem = DispatchDeliveryItem & {
unitIndex: number
sendType: number
note: string
responseMessage: string
}
export type DispatchStockItem = DispatchDeliveryItem & {
requiredCount: number
knapsackCount: number
purchasedCount: number
}
export type DispatchStockResult = {
usedKnapsack: boolean
purchaseTriggered: boolean
assetBefore: number
assetAfter: number
items: DispatchStockItem[]
}
export function resolveDispatchDeliveryItems(flow: JsonObject): DispatchDeliveryItem[] {
const rawItems = Array.isArray(flow.deliveryItems) ? flow.deliveryItems : []
const items = rawItems
.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
if (!Number.isInteger(cloudSkuId) || cloudSkuId <= 0) {
return null
}
return {
cloudSkuId,
cloudSkuName: String(source.cloudSkuName || source.skuName || '').trim(),
quantity: Number.isInteger(quantity) && quantity > 0 ? quantity : 1,
}
})
.filter((item): item is DispatchDeliveryItem => Boolean(item))
if (items.length > 0) {
return items
}
const binding = flow.binding && typeof flow.binding === 'object'
? (flow.binding as JsonObject)
: {}
const fallbackSkuId = Number(binding.skuId || 0) || 0
if (!fallbackSkuId) {
return []
}
return [
{
cloudSkuId: fallbackSkuId,
cloudSkuName: String(binding.skuName || '').trim(),
quantity: 1,
},
]
}
@@ -1,6 +1,9 @@
/**
* kuaishou-lewancloud)履约入口。
* 实现已按职责拆分到 prepare / rebind / probe / refresh-role 等模块。
* kuaishou-lewan业务名)履约入口。
* 模块路径仍为 kuaishou-cloudexecutor_key = kuaishou_ct_assisted
* 外部平台 API 在 platforms/cloudtentacles。
*
* 实现按职责拆分:prepare / redeem / dispatch / return / role 等。
*/
export {
isKuaishouCloudBindUrlFresh,
@@ -27,7 +30,15 @@ export {
export { rebindKuaishouCloudTaskRole } from './rebind-role.js'
export { probeKuaishouCloudTaskBindUrl } from './probe-bind-url.js'
export { refreshKuaishouCloudTaskRoleInfo } from './refresh-role-info.js'
export { confirmKuaishouCloudTaskRole } from './confirm-role.js'
export type { ConfirmKuaishouCloudRoleResult } from './confirm-role.js'
export { dispatchKuaishouCloudFulfillmentTask } from './dispatch-fulfillment.js'
export { returnKuaishouCloudFulfillmentTask } from './return-fulfillment.js'
export { redeemKuaishouCloudTask } from './redeem-fulfillment.js'
export type { RedeemKuaishouCloudTaskResult } from './redeem-fulfillment.js'
export {
dispatchKuaishouCloudFulfillmentTask,
returnKuaishouCloudFulfillmentTask,
} from './task-finalization.js'
isKuaishouCloudMockTask,
isKuaishouCloudMockContext,
buildMockVerifiedKuaishouCloudFlow,
completeMockKuaishouCloudTask,
} from './mock-helpers.js'
@@ -0,0 +1,123 @@
import { createTaskEvent } from '../../../repositories/task-event-repo.js'
import { updateTask } from '../../../repositories/task-repo.js'
import { TASK_STATUS } from '../../../domain/task-status.js'
import { asJsonObject, type JsonObject } from '../../../types/json.js'
import type { TaskRow } from '../../../types/repository/rows.js'
import { addHours } from '../../../utils/time.js'
import { normalizeKuaishouCloudFlow } from './domain.js'
import { parseTaskContext } from './task-context.js'
export function isKuaishouCloudMockTask(task: Partial<TaskRow> | null | undefined) {
return isKuaishouCloudMockContext(parseTaskContext(task))
}
export function isKuaishouCloudMockContext(context: JsonObject = {}) {
const flow = asJsonObject(context.kuaishouCloudFulfillment)
const mock = flow.mock && typeof flow.mock === 'object' ? flow.mock : context.mock
return Boolean(mock && typeof mock === 'object' && (mock as JsonObject).enabled === true)
}
export function buildMockVerifiedKuaishouCloudFlow(value: unknown, timestamp: string) {
const flow = normalizeKuaishouCloudFlow(value)
const source = flow as JsonObject
const bindUrl =
flow.binding.bindUrl ||
`https://example.com/mock-kuaishou-cloud-bind?task=mock&ts=${encodeURIComponent(timestamp)}`
const vnPhone = flow.binding.vnPhone || '13800000000'
const roleName = flow.binding.roleName || flow.role.name || '测试角色'
const roleId = flow.binding.roleId || flow.role.rid || '10001'
return {
...flow,
mock: {
...asJsonObject(source.mock),
enabled: true,
},
binding: {
...flow.binding,
prepareStatus: 'ready',
vnId: flow.binding.vnId || 900001,
vnPhone,
bindUrl,
bindPreparedAt: flow.binding.bindPreparedAt || timestamp,
bindExpiresAt: flow.binding.bindExpiresAt || addHours(timestamp, 24),
roleName,
roleId,
},
role: {
...flow.role,
status: 'ready',
name: roleName,
rid: roleId,
refreshedAt: flow.role.refreshedAt || timestamp,
errorMessage: '',
rawInfo: flow.role.rawInfo || {
mock: true,
},
},
}
}
/** 开发 mock:跳过真实 CT,直接落成 completed */
export async function completeMockKuaishouCloudTask(task: TaskRow, timestamp: string) {
const taskContext = parseTaskContext(task)
const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment)
const source = flow as JsonObject
const nextFlow = {
...flow,
mock: {
...asJsonObject(source.mock),
enabled: true,
},
dispatch: {
...flow.dispatch,
status: 'success',
dispatchAt: timestamp,
dispatchBy: {
source: 'claim_page_mock',
},
note: '开发 mock 已模拟发货成功',
items: flow.deliveryItems,
},
returnNumber: {
...flow.returnNumber,
status: 'success',
returnedAt: timestamp,
returnedBy: {
source: 'claim_page_mock',
},
},
consume: {
...flow.consume,
status: 'success',
consumedAt: timestamp,
errorMessage: '',
},
}
await updateTask(task.id, {
task_status: TASK_STATUS.COMPLETED,
delivery_status: 'success',
result_code: 'mock_success',
result_message: '开发 mock 已模拟兑换成功',
user_action_status: 'not_required',
last_error: '',
context_json: JSON.stringify({
...taskContext,
kuaishouCloudFulfillment: nextFlow,
}),
redeemed_at: timestamp,
updated_at: timestamp,
})
await createTaskEvent(
task.id,
'kuaishou_cloud_mock_redeemed',
{
source: 'claim_page_mock',
deliveryItems: flow.deliveryItems,
},
timestamp,
)
}
@@ -0,0 +1,95 @@
import assert from 'node:assert/strict'
import test from 'node:test'
import {
TASK_STATUS,
canRedeemKuaishouCloudClaimStatus,
isKuaishouCloudRedeemSettledStatus,
} from '../../../domain/task-status.js'
import { assertBoundUidMatchesExpected, assertClaimExpectedUidReady } from '../../claim/claim-identity.js'
/**
* redeem 用例前置条件单测(不打 DB / CT)。
* 完整 dispatch 走集成路径;这里锁定状态机与 UID 闸门契约。
*/
function canEnterRedeem(options: {
status: unknown
expectedUid: string
boundUid: string
mock?: boolean
}) {
if (isKuaishouCloudRedeemSettledStatus(options.status)) {
return { ok: true as const, alreadySettled: true }
}
if (!canRedeemKuaishouCloudClaimStatus(options.status)) {
return { ok: false as const, reason: 'status_not_redeemable' }
}
if (!options.mock && !options.expectedUid) {
return { ok: false as const, reason: 'uid_missing' }
}
if (!options.mock && options.expectedUid !== options.boundUid) {
return { ok: false as const, reason: 'uid_mismatch' }
}
return { ok: true as const, alreadySettled: false }
}
test('redeem 前置:waiting_binding + uid 匹配可进入', () => {
const result = canEnterRedeem({
status: TASK_STATUS.WAITING_BINDING,
expectedUid: '10001',
boundUid: '10001',
})
assert.deepEqual(result, { ok: true, alreadySettled: false })
})
test('redeem 前置:role_confirmed + uid 匹配可进入', () => {
const result = canEnterRedeem({
status: TASK_STATUS.ROLE_CONFIRMED,
expectedUid: '10001',
boundUid: '10001',
})
assert.deepEqual(result, { ok: true, alreadySettled: false })
})
test('redeem 前置:已 completed 视为 settled', () => {
const result = canEnterRedeem({
status: TASK_STATUS.COMPLETED,
expectedUid: '10001',
boundUid: '10001',
})
assert.deepEqual(result, { ok: true, alreadySettled: true })
})
test('redeem 前置:uid 不一致拒绝', () => {
const result = canEnterRedeem({
status: TASK_STATUS.WAITING_BINDING,
expectedUid: '10001',
boundUid: '99999',
})
assert.equal(result.ok, false)
if (!result.ok) {
assert.equal(result.reason, 'uid_mismatch')
}
})
test('redeem 闸门函数与 claim-identity 一致', () => {
assert.throws(
() => assertClaimExpectedUidReady({ context_json: '{}' }),
/填写游戏 UID/,
)
assert.throws(
() =>
assertBoundUidMatchesExpected(
{
context_json: JSON.stringify({
claimIdentity: { expectedUid: '10001' },
}),
},
{
binding: { roleId: 'x', roleName: 'n', vnPhone: '1' },
role: { rid: 'x', name: 'n' },
},
),
/不一致/,
)
})
@@ -0,0 +1,200 @@
import { createTaskEvent } from '../../../repositories/task-event-repo.js'
import {
getTaskById,
updateTask,
updateTaskStatusIfCurrent,
} from '../../../repositories/task-repo.js'
import {
TASK_STATUS,
canRedeemKuaishouCloudClaimStatus,
isKuaishouCloudRedeemSettledStatus,
normalizeTaskStatus,
} from '../../../domain/task-status.js'
import { createHttpError } from '../../../utils/http.js'
import { nowIso } from '../../../utils/time.js'
import type { TaskRow } from '../../../types/repository/rows.js'
import {
assertBoundUidMatchesExpected,
assertClaimExpectedUidReady,
getClaimIdentityFromTask,
} from '../../claim/claim-identity.js'
import { isKuaishouCloudTask, normalizeKuaishouCloudFlow, type JsonObject } from './domain.js'
import { dispatchKuaishouCloudFulfillmentTask } from './dispatch-fulfillment.js'
import {
completeMockKuaishouCloudTask,
isKuaishouCloudMockTask,
} from './mock-helpers.js'
import { normalizeActor, parseTaskContext } from './task-context.js'
export type RedeemKuaishouCloudTaskResult = {
task: TaskRow
/** 是否已处于终态/已兑换,调用方无需再 dispatch */
alreadySettled: boolean
/** 并发抢锁失败(他人已推进) */
concurrentSkipped: boolean
}
/**
* lewan 一键兑换用例(履约引擎入口):
* waiting_binding + UID 匹配 → role_confirmed → redeeming → dispatch(+autoFinalize)
*
* claim / admin 应只调此函数,不要直接拼状态机。
*/
export async function redeemKuaishouCloudTask(
task: TaskRow,
options: JsonObject = {},
): Promise<RedeemKuaishouCloudTaskResult> {
if (!isKuaishouCloudTask(task)) {
throw createHttpError('当前任务不是 kuaishou-lewan 履约任务', {
statusCode: 409,
errorCode: 'kuaishou_cloud_task_invalid',
})
}
const now = String(options.now || '').trim() || nowIso()
const source = String(options.source || 'system_redeem').trim() || 'system_redeem'
const actor = normalizeActor(options.actor) || { source }
const autoFinalize = options.autoFinalize !== false
const errorCodePrefix =
String(options.errorCodePrefix || 'kuaishou_cloud_redeem').trim() || 'kuaishou_cloud_redeem'
let workingTask = task
const currentStatus = normalizeTaskStatus(workingTask.task_status)
if (isKuaishouCloudRedeemSettledStatus(currentStatus)) {
return { task: workingTask, alreadySettled: true, concurrentSkipped: false }
}
if (!canRedeemKuaishouCloudClaimStatus(currentStatus)) {
throw createHttpError('当前状态不可兑换,请先完成绑定并匹配 UID', {
statusCode: 409,
errorCode: `${errorCodePrefix}_not_ready`,
})
}
// WAITING_BINDING + UID 匹配:自动升为 ROLE_CONFIRMED,实现一键兑换
if (currentStatus === TASK_STATUS.WAITING_BINDING) {
const flow = normalizeKuaishouCloudFlow(
parseTaskContext(workingTask).kuaishouCloudFulfillment,
)
if (!isKuaishouCloudMockTask(workingTask)) {
assertBoundUidMatchesExpected(workingTask, flow, {
errorCodePrefix,
})
} else {
assertClaimExpectedUidReady(workingTask)
}
const confirmed = await updateTask(workingTask.id, {
task_status: TASK_STATUS.ROLE_CONFIRMED,
role_id: flow.binding.roleId || workingTask.role_id || '',
role_name: flow.binding.roleName || workingTask.role_name || '',
user_action_status: TASK_STATUS.ROLE_CONFIRMED,
role_confirmed_at: now,
last_error: '',
updated_at: now,
})
if (confirmed) {
workingTask = confirmed
}
await createTaskEvent(
workingTask.id,
'kuaishou_cloud_role_confirmed',
{
expectedUid: getClaimIdentityFromTask(workingTask).expectedUid,
vnPhone: flow.binding.vnPhone,
roleName: flow.binding.roleName,
roleId: flow.binding.roleId,
source: `${source}_one_click`,
},
now,
)
}
const flowBeforeRedeem = normalizeKuaishouCloudFlow(
parseTaskContext(workingTask).kuaishouCloudFulfillment,
)
if (!isKuaishouCloudMockTask(workingTask)) {
assertBoundUidMatchesExpected(workingTask, flowBeforeRedeem, {
errorCodePrefix,
})
} else {
assertClaimExpectedUidReady(workingTask)
}
const lockedTask = await updateTaskStatusIfCurrent(
workingTask.id,
TASK_STATUS.ROLE_CONFIRMED,
{
task_status: TASK_STATUS.REDEEMING,
user_action_status: 'not_required',
last_error: '',
updated_at: now,
},
)
if (!lockedTask) {
const latest = (await getTaskById(workingTask.id)) || workingTask
return {
task: latest,
alreadySettled: isKuaishouCloudRedeemSettledStatus(latest.task_status),
concurrentSkipped: true,
}
}
if (isKuaishouCloudMockTask(lockedTask)) {
await completeMockKuaishouCloudTask(lockedTask, now)
const completed = (await getTaskById(lockedTask.id)) || lockedTask
return { task: completed, alreadySettled: false, concurrentSkipped: false }
}
try {
const result = await dispatchKuaishouCloudFulfillmentTask(lockedTask, {
source,
actor,
autoFinalize,
})
return {
task: result.task || lockedTask,
alreadySettled: false,
concurrentSkipped: false,
}
} catch (error) {
const latestTask = await getTaskById(lockedTask.id)
const latestStatus = latestTask ? normalizeTaskStatus(latestTask.task_status) : ''
if (latestTask && latestStatus !== TASK_STATUS.REDEEMING) {
if (
latestStatus === TASK_STATUS.DISPATCHED_PENDING_RETURN ||
latestStatus === TASK_STATUS.COMPLETED ||
latestStatus === TASK_STATUS.REDEEMED
) {
return { task: latestTask, alreadySettled: true, concurrentSkipped: false }
}
throw error
}
const message =
error instanceof Error ? error.message : '兑换请求提交失败,请联系客服处理'
await updateTask(lockedTask.id, {
task_status: TASK_STATUS.MANUAL_REVIEW,
user_action_status: 'not_required',
last_error: message,
result_code: 'kuaishou_cloud_redeem_failed',
result_message: message,
updated_at: nowIso(),
})
await createTaskEvent(
lockedTask.id,
'kuaishou_cloud_redeem_failed',
{
source,
errorMessage: message,
},
nowIso(),
)
throw error
}
}
@@ -0,0 +1,251 @@
import { createTaskEvent } from '../../../repositories/task-event-repo.js'
import { updateTask } from '../../../repositories/task-repo.js'
import { createHttpError } from '../../../utils/http.js'
import { nowIso } from '../../../utils/time.js'
import { TASK_STATUS, normalizeTaskStatus, type TaskStatus } from '../../../domain/task-status.js'
import { backCloudtentaclesVirtualNumber } from '../../platforms/cloudtentacles/virtual-number-service.js'
import { consumeKuaishouIndustryVouchersForTask } from '../../platforms/kuaishou-industry/voucher-service.js'
import type { TaskRow } from '../../../types/repository/rows.js'
import {
isIndustryEVoucherTask,
isKuaishouCloudTask,
maskCode,
maskPhone,
normalizeKuaishouCloudFlow,
type JsonObject,
} from './domain.js'
import { resolvePersistedCloudtentaclesContextBySourceKeys } from './cloudtentacles-context.js'
import { normalizeActor, parseTaskContext } from './task-context.js'
/**
* lewan 退还虚拟号。
*
* - 默认(发货后 autoFinalize):退号 + 尝试行业电子凭证核销收尾
* - `consumeIndustryVoucher: false`admin 清理):只退虚拟号,避免占号;与核销无关
*/
export async function returnKuaishouCloudFulfillmentTask(
task: TaskRow,
options: JsonObject = {},
) {
if (!isKuaishouCloudTask(task)) {
throw createHttpError('当前任务不是 kuaishou-lewan 履约任务', {
statusCode: 409,
errorCode: 'kuaishou_cloud_task_invalid',
})
}
const actor = normalizeActor(options.actor)
const now = nowIso()
const taskContext = parseTaskContext(task)
const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment)
const cloudContext = resolvePersistedCloudtentaclesContextBySourceKeys([
flow.binding.resolvedSourceKey,
...flow.binding.cloudSourceKeys,
])
// admin 退号是清理虚拟号;只有发货收尾路径才顺带核销
const allowConsumeIndustryVoucher = options.consumeIndustryVoucher !== false
if (!flow.binding.vnId || !flow.binding.vnKey) {
throw createHttpError('当前任务缺少可退还的虚拟号信息', {
statusCode: 409,
errorCode: 'kuaishou_cloud_missing_return_context',
})
}
const taskStatus = normalizeTaskStatus(task.task_status)
if (
flow.returnNumber.status === 'success' &&
(taskStatus === TASK_STATUS.COMPLETED || taskStatus === TASK_STATUS.MANUAL_REVIEW)
) {
return { task, flow }
}
await backCloudtentaclesVirtualNumber({
...cloudContext,
key: flow.binding.vnKey,
id: flow.binding.vnId,
})
const ticketCode = String(flow.ticket.code || '').trim()
let consumeStatus = 'pending'
let consumeErrorMessage = ''
let consumedAt: string | null = null
let nextTaskStatus: TaskStatus = TASK_STATUS.COMPLETED
let nextResultCode = 'kuaishou_cloud_completed'
const consumeAlreadyCompleted = flow.consume.status === 'success'
const hasIndustryVoucher = hasKuaishouIndustryVoucherContext(taskContext)
const isIndustryTask = isIndustryEVoucherTask(task)
const shouldConsumeIndustryVoucher =
allowConsumeIndustryVoucher && (hasIndustryVoucher || isIndustryTask)
let nextResultMessage = shouldConsumeIndustryVoucher
? 'cloudtentacles 发货、退号并完成电子凭证核销'
: allowConsumeIndustryVoucher
? 'cloudtentacles 发货、退号并收口'
: '虚拟号已退还(清理占号,未执行核销)'
let industryVoucherContextPatch: JsonObject | null = null
if (consumeAlreadyCompleted) {
consumeStatus = 'success'
consumedAt = (flow.consume.consumedAt as string | null) || now
if (!allowConsumeIndustryVoucher) {
nextResultMessage = '虚拟号已退还;电子凭证此前已核销'
}
} else if (shouldConsumeIndustryVoucher) {
const industryResult = hasIndustryVoucher
? await consumeKuaishouIndustryVouchersForTask(task, {
source:
String(options.source || 'system_auto_finalize').trim() || 'system_auto_finalize',
token: String(
isPlainObject(taskContext.kuaishouIndustryVoucher)
? taskContext.kuaishouIndustryVoucher.token || ''
: '',
).trim(),
consumeTime: Date.now(),
})
: { ok: true, consumed: [] as Array<Record<string, unknown>>, failed: [] as Array<{ errorMessage?: string }> }
if (industryResult.ok) {
consumeStatus = 'success'
consumedAt = now
const consumedVoucher = industryResult.consumed[0] || null
if (consumedVoucher) {
industryVoucherContextPatch = {
oid: consumedVoucher.oid,
token: consumedVoucher.token,
eticketId: consumedVoucher.voucher_code,
voucherCode: consumedVoucher.voucher_code,
unitIndex: Number(consumedVoucher.unit_index || 0) || 0,
status: 'CONSUMED',
validStartTime: Number(consumedVoucher.valid_start_time || 0) || 0,
validEndTime: Number(consumedVoucher.valid_end_time || 0) || 0,
consumedAt: now,
consumeSerialNum:
consumedVoucher.consume_serial_num || `CONSUME-${consumedVoucher.voucher_code}`,
}
}
} else {
consumeStatus = 'failed'
consumeErrorMessage =
industryResult.failed[0]?.errorMessage || '电子凭证核销回调失败,请人工处理'
}
} else {
// cleanup 或无凭证:不碰核销
consumeStatus = 'skipped'
consumeErrorMessage = ''
}
if (consumeStatus !== 'success') {
if (shouldConsumeIndustryVoucher) {
nextTaskStatus = TASK_STATUS.MANUAL_REVIEW
nextResultCode = 'kuaishou_cloud_consume_failed'
nextResultMessage =
consumeErrorMessage || '号码已退还,但电子凭证核销未完成,请人工处理'
} else {
nextResultCode = allowConsumeIndustryVoucher
? 'kuaishou_cloud_completed_without_eticket_consume'
: 'kuaishou_cloud_number_returned_cleanup'
}
}
const nextContext = {
...taskContext,
...(industryVoucherContextPatch
? {
kuaishouIndustryVoucher: {
...(isPlainObject(taskContext.kuaishouIndustryVoucher)
? taskContext.kuaishouIndustryVoucher
: {}),
...industryVoucherContextPatch,
},
}
: {}),
kuaishouCloudFulfillment: {
...flow,
returnNumber: {
...flow.returnNumber,
status: 'success',
returnedAt: now,
returnedBy: actor,
},
consume: {
...flow.consume,
status: consumeStatus,
shopId: flow.consume.shopId,
shopName: flow.consume.shopName,
autoConsumeEnabled: flow.consume.autoConsumeEnabled === true,
consumedAt,
errorMessage: consumeErrorMessage,
},
},
}
const updatedTask = await updateTask(task.id, {
task_status: nextTaskStatus,
delivery_status: 'delivered',
result_code: nextResultCode,
result_message: nextResultMessage,
redeemed_at: consumeStatus === 'failed' ? task.redeemed_at : now,
last_error: consumeErrorMessage,
context_json: JSON.stringify(nextContext),
updated_at: now,
})
await createTaskEvent(
task.id,
'kuaishou_cloud_number_returned',
{
source: String(options.source || 'system').trim() || 'system',
vnId: flow.binding.vnId,
vnPhoneMasked: maskPhone(flow.binding.vnPhone),
actor,
},
now,
)
const consumeEventType = consumeAlreadyCompleted
? 'kuaishou_cloud_consume_already_completed'
: consumeStatus === 'success'
? 'kuaishou_cloud_consumed'
: consumeStatus === 'skipped'
? 'kuaishou_cloud_consume_skipped'
: 'kuaishou_cloud_consume_failed'
await createTaskEvent(
task.id,
consumeEventType,
{
source: String(options.source || 'system').trim() || 'system',
ticketCodeMasked: maskCode(ticketCode),
shopId: flow.consume.shopId,
shopName: flow.consume.shopName,
consumeStatus,
consumeMode: shouldConsumeIndustryVoucher
? 'industry_voucher'
: allowConsumeIndustryVoucher
? 'legacy_writeoff_disabled'
: 'cleanup_skip_consume',
errorMessage: consumeErrorMessage,
actor,
},
now,
)
return {
task: updatedTask,
flow: normalizeKuaishouCloudFlow(parseTaskContext(updatedTask).kuaishouCloudFulfillment),
}
}
function hasKuaishouIndustryVoucherContext(value: JsonObject): boolean {
const voucher = isPlainObject(value.kuaishouIndustryVoucher)
? value.kuaishouIndustryVoucher
: {}
const voucherCode = String(voucher.voucherCode || voucher.eticketId || '').trim()
const status = String(voucher.status || 'UNUSED').trim().toUpperCase()
const sendCallbackStatus = String(voucher.sendCallbackStatus || 'success').trim().toLowerCase()
return Boolean(voucherCode && status !== 'DESTROYED' && sendCallbackStatus === 'success')
}
function isPlainObject(value: unknown): value is JsonObject {
return Boolean(value) && typeof value === 'object' && !Array.isArray(value)
}
@@ -1,940 +1,12 @@
import { createTaskEvent } from "../../../repositories/task-event-repo.js";
import { updateTask } from "../../../repositories/task-repo.js";
import { createHttpError } from "../../../utils/http.js";
import { nowIso } from "../../../utils/time.js";
import { TASK_STATUS, normalizeTaskStatus, type TaskStatus } from "../../../domain/task-status.js";
import { notifyKuaishouCloudAssetNotEnough } from "../../notification/domain-notifications.js";
import {
buyCloudtentaclesSku,
getCloudtentaclesAsset,
listCloudtentaclesSku,
useCloudtentaclesSku,
} from "../../platforms/cloudtentacles/catalog-service.js";
import { getCloudtentaclesKnapsack } from "../../platforms/cloudtentacles/knapsack-service.js";
import { backCloudtentaclesVirtualNumber } from "../../platforms/cloudtentacles/virtual-number-service.js";
import { consumeKuaishouIndustryVouchersForTask } from "../../platforms/kuaishou-industry/voucher-service.js";
import {
isKuaishouCloudTask,
isIndustryEVoucherTask,
maskCode,
maskPhone,
normalizeKuaishouCloudFlow,
type JsonObject,
} from "./domain.js";
import { resolvePersistedCloudtentaclesContextBySourceKeys } from "./cloudtentacles-context.js";
import { syncKuaishouCloudRoleInfoBeforeDispatch } from "./dispatch-role-sync.js";
import { normalizeActor, parseTaskContext } from "./task-context.js";
import type { TaskRow } from "../../../types/repository/rows.js";
type DispatchDeliveryItem = {
cloudSkuId: number;
cloudSkuName: string;
quantity: number;
};
type DispatchResultItem = DispatchDeliveryItem & {
unitIndex: number;
sendType: number;
note: string;
responseMessage: string;
};
type DispatchStockItem = DispatchDeliveryItem & {
requiredCount: number;
knapsackCount: number;
purchasedCount: number;
};
type DispatchStockResult = {
usedKnapsack: boolean;
purchaseTriggered: boolean;
assetBefore: number;
assetAfter: number;
items: DispatchStockItem[];
};
const cloudSkuDispatchLocks = new Map<string, Promise<void>>();
export async function dispatchKuaishouCloudFulfillmentTask(
task: TaskRow,
options: JsonObject = {}
) {
if (!isKuaishouCloudTask(task)) {
throw createHttpError("当前任务不是 kuaishou-lewan 履约任务", {
statusCode: 409,
errorCode: "kuaishou_cloud_task_invalid",
});
}
const actor = normalizeActor(options.actor);
const now = nowIso();
const taskContext = parseTaskContext(task);
const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment);
const cloudContext = resolvePersistedCloudtentaclesContextBySourceKeys([
flow.binding.resolvedSourceKey,
...flow.binding.cloudSourceKeys,
]);
if (!flow.binding.skuId || !flow.binding.vnId || !flow.binding.vnPhone) {
throw createHttpError(
"当前任务还没有准备好绑定资源,请先完成绑定资源准备",
{
statusCode: 409,
errorCode: "kuaishou_cloud_not_prepared",
}
);
}
const persistedTicketCode = String(flow.ticket.code || "").trim();
const voucherContext = isPlainObject(taskContext.kuaishouIndustryVoucher)
? taskContext.kuaishouIndustryVoucher
: {};
const industryVoucherCode = String(
voucherContext.voucherCode || voucherContext.eticketId || ""
).trim();
const hasIndustryVoucherForDispatch =
Boolean(industryVoucherCode) || isIndustryEVoucherTask(task);
const resolvedTicketCode = persistedTicketCode || industryVoucherCode;
if (!resolvedTicketCode && !hasIndustryVoucherForDispatch) {
throw createHttpError("旧快手小店核销流程已停用,请改用行业电子凭证处理", {
statusCode: 409,
errorCode: "kuaishou_cloud_missing_ticket_code",
});
}
if (
flow.dispatch.status === "success" &&
normalizeTaskStatus(task.task_status) === TASK_STATUS.DISPATCHED_PENDING_RETURN
) {
return { task, flow };
}
const synced = await syncKuaishouCloudRoleInfoBeforeDispatch(task, {
now,
actor,
source: options.source || "system_before_dispatch",
cloudContext,
taskContext,
flow,
});
const syncedFlow = synced.flow;
const syncedTaskContext = synced.taskContext;
const deliveryItems = resolveDispatchDeliveryItems(syncedFlow);
if (deliveryItems.length === 0) {
throw createHttpError("当前任务缺少 cloud 发货物品配置", {
statusCode: 409,
errorCode: "kuaishou_cloud_missing_delivery_items",
});
}
let dispatchResults: DispatchResultItem[];
let stockResult: DispatchStockResult;
try {
({ dispatchResults, stockResult } = await withCloudSkuDispatchLocks(
resolveCloudSkuDispatchLockKeys(cloudContext.resolvedSourceKey, deliveryItems),
() => prepareStockAndDispatch({
task,
flow: syncedFlow,
cloudContext,
deliveryItems,
})
));
} catch (error) {
await markKuaishouCloudDispatchFailed(task, syncedTaskContext, syncedFlow, error, {
actor,
source: options.source || "system",
});
throw error;
}
const firstDispatchResult: DispatchResultItem = dispatchResults[0] || {
cloudSkuId: syncedFlow.binding.skuId,
cloudSkuName: syncedFlow.binding.skuName,
quantity: 1,
unitIndex: 1,
sendType: 0,
note: "",
responseMessage: "",
};
const dispatchSummary =
dispatchResults.length > 1
? `cloudtentacles 发货成功,共 ${dispatchResults.length}`
: String(
firstDispatchResult.responseMessage ||
firstDispatchResult.note ||
"cloudtentacles 发货成功"
).trim();
const nextContext = {
...syncedTaskContext,
kuaishouCloudFulfillment: {
...syncedFlow,
ticket: {
...syncedFlow.ticket,
code: resolvedTicketCode,
capturedAt: resolvedTicketCode
? syncedFlow.ticket.capturedAt || now
: syncedFlow.ticket.capturedAt,
capturedBy: syncedFlow.ticket.capturedBy,
},
dispatch: {
...syncedFlow.dispatch,
status: "success",
dispatchAt: now,
dispatchBy: actor,
sendType: Number(firstDispatchResult.sendType || 0) || 0,
note: dispatchSummary,
items: dispatchResults,
},
purchase: {
...syncedFlow.purchase,
usedKnapsack: stockResult.usedKnapsack,
purchaseTriggered: stockResult.purchaseTriggered,
assetBefore: stockResult.assetBefore,
assetAfter: stockResult.assetAfter,
purchaseAt: stockResult.purchaseTriggered
? now
: syncedFlow.purchase.purchaseAt,
items: stockResult.items,
},
},
};
let updatedTask = await updateTask(task.id, {
task_status: TASK_STATUS.DISPATCHED_PENDING_RETURN,
delivery_status: "delivered",
result_code: "kuaishou_cloud_dispatched",
result_message: dispatchSummary,
user_action_status: "not_required",
last_error: "",
context_json: JSON.stringify(nextContext),
updated_at: now,
});
if (!updatedTask) {
throw createHttpError("kuaishou-lewan 发货状态更新失败", {
statusCode: 500,
errorCode: "kuaishou_cloud_dispatch_update_failed",
});
}
await createTaskEvent(
task.id,
"kuaishou_cloud_dispatched",
{
source: String(options.source || "system").trim() || "system",
ticketCodeMasked: maskCode(resolvedTicketCode),
skuId: syncedFlow.binding.skuId,
vnId: syncedFlow.binding.vnId,
vnPhoneMasked: maskPhone(syncedFlow.binding.vnPhone),
sendType: firstDispatchResult.sendType,
note: dispatchSummary,
deliveryItems,
dispatchResults,
stockResult,
actor,
},
now
);
const shouldAutoFinalize =
options.autoFinalize === true &&
normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment)
.returnNumber.autoReturnEnabled === true;
if (shouldAutoFinalize) {
const finalizeResult = await returnKuaishouCloudFulfillmentTask(
updatedTask,
{
actor,
source: options.source || "system_auto_finalize",
}
);
updatedTask = finalizeResult.task;
}
return {
task: updatedTask,
flow: normalizeKuaishouCloudFlow(
parseTaskContext(updatedTask).kuaishouCloudFulfillment
),
};
}
async function markKuaishouCloudDispatchFailed(
task: TaskRow,
taskContext: JsonObject,
flow: JsonObject,
error: unknown,
options: JsonObject = {}
) {
const now = nowIso();
const errorMessage = resolveErrorMessage(error);
const errorCode = resolveErrorCode(error) || "kuaishou_cloud_dispatch_failed";
const failureContext = resolveDispatchFailureContext(error);
const stockResult = isPlainObject(failureContext.stockResult)
? (failureContext.stockResult as Partial<DispatchStockResult>)
: null;
const nextContext = {
...taskContext,
kuaishouCloudFulfillment: {
...flow,
dispatch: {
...flow.dispatch,
status: "failed",
failedAt: now,
failedStage: String(failureContext.stage || "").trim(),
errorCode,
errorMessage,
note: buildDispatchFailedMessage(errorMessage),
},
purchase: {
...flow.purchase,
...(stockResult
? {
usedKnapsack: stockResult.usedKnapsack,
purchaseTriggered: stockResult.purchaseTriggered,
assetBefore: stockResult.assetBefore,
assetAfter: stockResult.assetAfter,
items: stockResult.items,
}
: {}),
},
},
};
await updateTask(task.id, {
task_status: TASK_STATUS.MANUAL_REVIEW,
user_action_status: "not_required",
last_error: errorMessage,
result_code: errorCode,
result_message: buildDispatchFailedMessage(errorMessage),
context_json: JSON.stringify(nextContext),
updated_at: now,
});
await createTaskEvent(
task.id,
"kuaishou_cloud_dispatch_failed",
{
source: String(options.source || "system").trim() || "system",
actor: options.actor,
errorCode,
errorMessage,
failureContext,
},
now
);
}
async function prepareStockAndDispatch({
task,
flow,
cloudContext,
deliveryItems,
}: {
task: TaskRow;
flow: JsonObject;
cloudContext: JsonObject;
deliveryItems: DispatchDeliveryItem[];
}) {
const [knapsack, skuList] = await Promise.all([
getCloudtentaclesKnapsack(cloudContext),
listCloudtentaclesSku(cloudContext),
]);
const skuItems = Array.isArray(skuList.items) ? skuList.items : [];
const knapsackItems = Array.isArray(knapsack.items) ? knapsack.items : [];
const stockItems = buildDispatchStockItems(deliveryItems, {
skuItems,
knapsackItems,
});
const missingItems = stockItems.filter((item) => item.purchasedCount > 0);
let assetBefore = 0;
let assetAfter = 0;
let purchaseTriggered = false;
if (missingItems.length > 0) {
if (flow.purchase?.autoBuyEnabled === false) {
throw createHttpError("背包中没有现成库存,且当前配置未开启自动购买", {
statusCode: 409,
errorCode: "kuaishou_cloud_auto_buy_disabled",
});
}
const missingSku = missingItems.find((item) => !findCloudSkuItem(skuItems, item.cloudSkuId));
if (missingSku) {
throw createHttpError(`cloudtentacles 未找到 SKU ${missingSku.cloudSkuId}`, {
statusCode: 404,
errorCode: "kuaishou_cloud_sku_not_found",
});
}
const asset = await getCloudtentaclesAsset(cloudContext);
assetBefore = Number(asset.asset || 0) || 0;
const targetPrice = missingItems.reduce((sum, item) => {
const skuItem = findCloudSkuItem(skuItems, item.cloudSkuId);
return sum + item.purchasedCount * Number(skuItem?.price || 0);
}, 0);
const requiredAsset = targetPrice + Number(flow.purchase?.minAssetReserve || 0);
if (assetBefore < requiredAsset) {
await notifyKuaishouCloudAssetNotEnough({
task,
flow,
assetBefore,
requiredAsset,
skuName: missingItems
.map((item) => item.cloudSkuName || `SKU ${item.cloudSkuId}`)
.join("、"),
});
throw createHttpError(`余额不足,当前 ${assetBefore},至少需要 ${requiredAsset}`, {
statusCode: 409,
errorCode: "kuaishou_cloud_asset_not_enough",
});
}
for (const item of missingItems) {
await createTaskEvent(
task.id,
"cloudtentacles_sku_buy_started",
{
source: "kuaishou_cloud_dispatch",
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
count: item.purchasedCount,
assetBefore,
},
nowIso()
);
try {
const buyResult = await buyCloudtentaclesSku({
...cloudContext,
id: item.cloudSkuId,
count: item.purchasedCount,
});
await createTaskEvent(
task.id,
"cloudtentacles_sku_buy_succeeded",
{
source: "kuaishou_cloud_dispatch",
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
count: item.purchasedCount,
responseMessage: buyResult.responseMessage,
},
nowIso()
);
} catch (error) {
await createTaskEvent(
task.id,
"cloudtentacles_sku_buy_failed",
{
source: "kuaishou_cloud_dispatch",
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
count: item.purchasedCount,
errorCode: resolveErrorCode(error),
errorMessage: resolveErrorMessage(error),
},
nowIso()
);
throw enrichCloudtentaclesDispatchError(error, {
stage: "purchase",
stockResult: buildCurrentStockResult({
stockItems,
purchaseTriggered,
assetBefore,
assetAfter,
}),
item,
});
}
}
purchaseTriggered = true;
const assetResult = await getCloudtentaclesAsset(cloudContext);
assetAfter = Number(assetResult.asset || 0) || 0;
}
const dispatchResults: DispatchResultItem[] = [];
for (const item of deliveryItems) {
for (let index = 0; index < item.quantity; index += 1) {
await createTaskEvent(
task.id,
"cloudtentacles_sku_use_started",
{
source: "kuaishou_cloud_dispatch",
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
unitIndex: index + 1,
quantity: item.quantity,
vnId: flow.binding.vnId,
vnPhoneMasked: maskPhone(flow.binding.vnPhone),
},
nowIso()
);
let dispatchResult: Awaited<ReturnType<typeof useCloudtentaclesSku>>;
try {
dispatchResult = await useCloudtentaclesSku({
...cloudContext,
id: item.cloudSkuId,
virtualNumberId: flow.binding.vnId,
phone: flow.binding.vnPhone,
});
} catch (error) {
await createTaskEvent(
task.id,
"cloudtentacles_sku_use_failed",
{
source: "kuaishou_cloud_dispatch",
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
unitIndex: index + 1,
quantity: item.quantity,
vnId: flow.binding.vnId,
vnPhoneMasked: maskPhone(flow.binding.vnPhone),
errorCode: resolveErrorCode(error),
errorMessage: resolveErrorMessage(error),
},
nowIso()
);
throw enrichCloudtentaclesDispatchError(error, {
stage: "dispatch",
stockResult: buildCurrentStockResult({
stockItems,
purchaseTriggered,
assetBefore,
assetAfter,
}),
item,
unitIndex: index + 1,
vnId: flow.binding.vnId,
vnPhoneMasked: maskPhone(flow.binding.vnPhone),
});
}
await createTaskEvent(
task.id,
"cloudtentacles_sku_use_succeeded",
{
source: "kuaishou_cloud_dispatch",
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
unitIndex: index + 1,
quantity: item.quantity,
vnId: flow.binding.vnId,
vnPhoneMasked: maskPhone(flow.binding.vnPhone),
sendType: Number(dispatchResult.sendType || 0) || 0,
note: String(dispatchResult.note || "").trim(),
responseMessage: String(dispatchResult.responseMessage || "").trim(),
},
nowIso()
);
dispatchResults.push({
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
unitIndex: index + 1,
quantity: item.quantity,
sendType: Number(dispatchResult.sendType || 0) || 0,
note: String(dispatchResult.note || "").trim(),
responseMessage: String(dispatchResult.responseMessage || "").trim(),
});
}
}
return {
dispatchResults,
stockResult: {
usedKnapsack: stockItems.every((item) => item.purchasedCount <= 0),
purchaseTriggered,
assetBefore,
assetAfter,
items: stockItems,
},
};
}
function buildCurrentStockResult({
stockItems,
purchaseTriggered,
assetBefore,
assetAfter,
}: {
stockItems: DispatchStockItem[];
purchaseTriggered: boolean;
assetBefore: number;
assetAfter: number;
}): DispatchStockResult {
return {
usedKnapsack: stockItems.every((item) => item.purchasedCount <= 0),
purchaseTriggered,
assetBefore,
assetAfter,
items: stockItems,
};
}
function enrichCloudtentaclesDispatchError(error: unknown, context: JsonObject) {
const currentError = error instanceof Error
? (error as Error & { context?: unknown })
: createHttpError(resolveErrorMessage(error), {
statusCode: 500,
errorCode: "kuaishou_cloud_dispatch_failed",
}) as Error & { context?: unknown };
const currentContext = isPlainObject(currentError.context) ? currentError.context : {};
currentError.context = {
...currentContext,
dispatchFailure: {
...(isPlainObject((currentContext as JsonObject).dispatchFailure)
? (currentContext as JsonObject).dispatchFailure
: {}),
...context,
},
};
return currentError;
}
function resolveDispatchFailureContext(error: unknown): JsonObject {
if (!error || typeof error !== "object") {
return {};
}
const context = (error as { context?: unknown }).context;
if (!isPlainObject(context)) {
return {};
}
const dispatchFailure = context.dispatchFailure;
return isPlainObject(dispatchFailure) ? dispatchFailure : {};
}
function resolveErrorCode(error: unknown) {
if (!error || typeof error !== "object") {
return "";
}
return String((error as { errorCode?: unknown }).errorCode || "").trim();
}
function resolveErrorMessage(error: unknown) {
return error instanceof Error ? error.message : String(error || "CloudTentacles 发货失败");
}
function buildDispatchFailedMessage(message: unknown) {
const reason = String(message || "").trim() || "未知错误";
return `CloudTentacles 发货失败:${reason}`;
}
function isPlainObject(value: unknown): value is JsonObject {
return Boolean(value) && typeof value === "object" && !Array.isArray(value);
}
function hasKuaishouIndustryVoucherContext(value: JsonObject): boolean {
const voucher = isPlainObject(value.kuaishouIndustryVoucher)
? value.kuaishouIndustryVoucher
: {};
const voucherCode = String(voucher.voucherCode || voucher.eticketId || "").trim();
const status = String(voucher.status || "UNUSED").trim().toUpperCase();
const sendCallbackStatus = String(voucher.sendCallbackStatus || "success").trim().toLowerCase();
return Boolean(voucherCode && status !== "DESTROYED" && sendCallbackStatus === "success");
}
export function buildDispatchStockItems(
deliveryItems: DispatchDeliveryItem[],
{ skuItems = [], knapsackItems = [] }: { skuItems?: JsonObject[]; knapsackItems?: JsonObject[] } = {}
): DispatchStockItem[] {
return deliveryItems.map((item) => {
const knapsackItem = findCloudSkuItem(knapsackItems, item.cloudSkuId);
const skuItem = findCloudSkuItem(skuItems, item.cloudSkuId);
const knapsackCount = Math.max(0, Number(knapsackItem?.count || 0) || 0);
return {
...item,
cloudSkuName: item.cloudSkuName || String(skuItem?.name || knapsackItem?.name || "").trim(),
requiredCount: item.quantity,
knapsackCount,
purchasedCount: Math.max(0, item.quantity - knapsackCount),
};
});
}
function findCloudSkuItem(items: JsonObject[], cloudSkuId: number) {
return items.find((item) => Number(item.id || 0) === cloudSkuId) || null;
}
function resolveCloudSkuDispatchLockKeys(sourceKey: unknown, deliveryItems: DispatchDeliveryItem[]) {
const normalizedSourceKey = String(sourceKey || "default").trim() || "default";
return deliveryItems.map((item) => `${normalizedSourceKey}:${item.cloudSkuId}`);
}
async function withCloudSkuDispatchLocks<T>(keys: string[], callback: () => Promise<T>): Promise<T> {
const normalizedKeys = [...new Set(keys.map((key) => String(key || "").trim()).filter(Boolean))].sort();
async function run(index: number): Promise<T> {
const key = normalizedKeys[index];
if (!key) {
return callback();
}
return withCloudSkuDispatchLock(key, () => run(index + 1));
}
return run(0);
}
async function withCloudSkuDispatchLock<T>(key: string, callback: () => Promise<T>): Promise<T> {
const previous = cloudSkuDispatchLocks.get(key) || Promise.resolve();
let release: () => void = () => {};
const current = new Promise<void>((resolve) => {
release = resolve;
});
const queued = previous.catch(() => {}).then(() => current);
cloudSkuDispatchLocks.set(key, queued);
await previous.catch(() => {});
try {
return await callback();
} finally {
release();
if (cloudSkuDispatchLocks.get(key) === queued) {
cloudSkuDispatchLocks.delete(key);
}
}
}
function resolveDispatchDeliveryItems(flow: JsonObject): DispatchDeliveryItem[] {
const rawItems = Array.isArray(flow.deliveryItems) ? flow.deliveryItems : [];
const items = rawItems
.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;
if (!Number.isInteger(cloudSkuId) || cloudSkuId <= 0) {
return null;
}
return {
cloudSkuId,
cloudSkuName: String(source.cloudSkuName || source.skuName || "").trim(),
quantity: Number.isInteger(quantity) && quantity > 0 ? quantity : 1,
};
})
.filter((item): item is DispatchDeliveryItem => Boolean(item));
if (items.length > 0) {
return items;
}
const fallbackSkuId = Number(flow?.binding?.skuId || 0) || 0;
if (!fallbackSkuId) {
return [];
}
return [
{
cloudSkuId: fallbackSkuId,
cloudSkuName: String(flow?.binding?.skuName || "").trim(),
quantity: 1,
},
];
}
export async function returnKuaishouCloudFulfillmentTask(
task: TaskRow,
options: JsonObject = {}
) {
if (!isKuaishouCloudTask(task)) {
throw createHttpError("当前任务不是 kuaishou-lewan 履约任务", {
statusCode: 409,
errorCode: "kuaishou_cloud_task_invalid",
});
}
const actor = normalizeActor(options.actor);
const now = nowIso();
const taskContext = parseTaskContext(task);
const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment);
const cloudContext = resolvePersistedCloudtentaclesContextBySourceKeys([
flow.binding.resolvedSourceKey,
...flow.binding.cloudSourceKeys,
]);
if (!flow.binding.vnId || !flow.binding.vnKey) {
throw createHttpError("当前任务缺少可退还的虚拟号信息", {
statusCode: 409,
errorCode: "kuaishou_cloud_missing_return_context",
});
}
const taskStatus = normalizeTaskStatus(task.task_status);
if (
flow.returnNumber.status === "success" &&
(taskStatus === TASK_STATUS.COMPLETED || taskStatus === TASK_STATUS.MANUAL_REVIEW)
) {
return { task, flow };
}
await backCloudtentaclesVirtualNumber({
...cloudContext,
key: flow.binding.vnKey,
id: flow.binding.vnId,
});
const ticketCode = String(flow.ticket.code || "").trim();
let consumeStatus = "pending";
let consumeErrorMessage = "";
let consumedAt = null;
let nextTaskStatus: TaskStatus = TASK_STATUS.COMPLETED;
let nextResultCode = "kuaishou_cloud_completed";
const consumeAlreadyCompleted = flow.consume.status === "success";
const hasIndustryVoucher = hasKuaishouIndustryVoucherContext(taskContext);
const isIndustryTask = isIndustryEVoucherTask(task);
const shouldConsumeIndustryVoucher = hasIndustryVoucher || isIndustryTask;
let nextResultMessage = shouldConsumeIndustryVoucher
? "cloudtentacles 发货、退号并完成电子凭证核销"
: "cloudtentacles 发货、退号并收口";
let industryVoucherContextPatch: JsonObject | null = null;
if (consumeAlreadyCompleted) {
consumeStatus = "success";
consumedAt = flow.consume.consumedAt || now;
} else if (shouldConsumeIndustryVoucher) {
const industryResult = hasIndustryVoucher
? await consumeKuaishouIndustryVouchersForTask(task, {
source: String(options.source || "system_auto_finalize").trim() || "system_auto_finalize",
token: String(taskContext.kuaishouIndustryVoucher?.token || "").trim(),
consumeTime: Date.now(),
})
: { ok: true, consumed: [], failed: [] };
if (industryResult.ok) {
consumeStatus = "success";
consumedAt = now;
const consumedVoucher = industryResult.consumed[0] || null;
if (consumedVoucher) {
industryVoucherContextPatch = {
oid: consumedVoucher.oid,
token: consumedVoucher.token,
eticketId: consumedVoucher.voucher_code,
voucherCode: consumedVoucher.voucher_code,
unitIndex: Number(consumedVoucher.unit_index || 0) || 0,
status: "CONSUMED",
validStartTime: Number(consumedVoucher.valid_start_time || 0) || 0,
validEndTime: Number(consumedVoucher.valid_end_time || 0) || 0,
consumedAt: now,
consumeSerialNum: consumedVoucher.consume_serial_num || `CONSUME-${consumedVoucher.voucher_code}`,
};
}
} else {
consumeStatus = "failed";
consumeErrorMessage =
industryResult.failed[0]?.errorMessage || "电子凭证核销回调失败,请人工处理";
}
} else {
consumeStatus = "skipped";
consumeErrorMessage = "";
}
if (consumeStatus !== "success") {
if (shouldConsumeIndustryVoucher) {
nextTaskStatus = TASK_STATUS.MANUAL_REVIEW;
nextResultCode = "kuaishou_cloud_consume_failed";
nextResultMessage =
consumeErrorMessage || "号码已退还,但电子凭证核销未完成,请人工处理";
} else {
nextResultCode = "kuaishou_cloud_completed_without_eticket_consume";
}
}
const nextContext = {
...taskContext,
...(industryVoucherContextPatch
? {
kuaishouIndustryVoucher: {
...(isPlainObject(taskContext.kuaishouIndustryVoucher)
? taskContext.kuaishouIndustryVoucher
: {}),
...industryVoucherContextPatch,
},
}
: {}),
kuaishouCloudFulfillment: {
...flow,
returnNumber: {
...flow.returnNumber,
status: "success",
returnedAt: now,
returnedBy: actor,
},
consume: {
...flow.consume,
status: consumeStatus,
shopId: flow.consume.shopId,
shopName: flow.consume.shopName,
autoConsumeEnabled: flow.consume.autoConsumeEnabled === true,
consumedAt,
errorMessage: consumeErrorMessage,
},
},
};
const updatedTask = await updateTask(task.id, {
task_status: nextTaskStatus,
delivery_status: "delivered",
result_code: nextResultCode,
result_message: nextResultMessage,
redeemed_at: consumeStatus === "failed" ? task.redeemed_at : now,
last_error: consumeErrorMessage,
context_json: JSON.stringify(nextContext),
updated_at: now,
});
await createTaskEvent(
task.id,
"kuaishou_cloud_number_returned",
{
source: String(options.source || "system").trim() || "system",
vnId: flow.binding.vnId,
vnPhoneMasked: maskPhone(flow.binding.vnPhone),
actor,
},
now
);
const consumeEventType = consumeAlreadyCompleted
? "kuaishou_cloud_consume_already_completed"
: consumeStatus === "success"
? "kuaishou_cloud_consumed"
: consumeStatus === "skipped"
? "kuaishou_cloud_consume_skipped"
: "kuaishou_cloud_consume_failed";
await createTaskEvent(
task.id,
consumeEventType,
{
source: String(options.source || "system").trim() || "system",
ticketCodeMasked: maskCode(ticketCode),
shopId: flow.consume.shopId,
shopName: flow.consume.shopName,
consumeStatus,
consumeMode: shouldConsumeIndustryVoucher ? "industry_voucher" : "legacy_writeoff_disabled",
errorMessage: consumeErrorMessage,
actor,
},
now
);
return {
task: updatedTask,
flow: normalizeKuaishouCloudFlow(
parseTaskContext(updatedTask).kuaishouCloudFulfillment
),
};
}
/**
* 兼容入口:原 task-finalization 已拆为
* - dispatch-fulfillment(发货编排)
* - dispatch-stock(库存/采购/下发)
* - dispatch-failure(失败落库)
* - return-fulfillment(退号 + 核销收尾)
*
* 新代码请优先从 index.js 或对应子模块 import。
*/
export { dispatchKuaishouCloudFulfillmentTask } from './dispatch-fulfillment.js'
export { returnKuaishouCloudFulfillmentTask } from './return-fulfillment.js'
export { buildDispatchStockItems } from './dispatch-stock.js'