diff --git a/apps/backend/src/services/admin/write/kuaishou-cloud-actions.ts b/apps/backend/src/services/admin/write/kuaishou-cloud-actions.ts index 35ac0aab..12775d88 100644 --- a/apps/backend/src/services/admin/write/kuaishou-cloud-actions.ts +++ b/apps/backend/src/services/admin/write/kuaishou-cloud-actions.ts @@ -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 { 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 { 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), } } diff --git a/apps/backend/src/services/claim/kuaishou-cloud-claim-service.ts b/apps/backend/src/services/claim/kuaishou-cloud-claim-service.ts index cda61e20..c222550a 100644 --- a/apps/backend/src/services/claim/kuaishou-cloud-claim-service.ts +++ b/apps/backend/src/services/claim/kuaishou-cloud-claim-service.ts @@ -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 +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( + result: Promise | T | null, + errorMessage: string, + errorCode: string, +): Promise> { + const value = await result + if (!value) { + throw createHttpError(errorMessage, { + statusCode: 409, + errorCode, + }) + } + return value +} + async function verifyIndustryVoucherTicket( context: Awaited>, 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 | 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, - ) -} diff --git a/apps/backend/src/services/claim/kuaishou-cloud-sync-service.ts b/apps/backend/src/services/claim/kuaishou-cloud-sync-service.ts index 5b81ad90..f51ddffb 100644 --- a/apps/backend/src/services/claim/kuaishou-cloud-sync-service.ts +++ b/apps/backend/src/services/claim/kuaishou-cloud-sync-service.ts @@ -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 } diff --git a/apps/backend/src/services/fulfillment/executors/kuaishou-cloud-executor.ts b/apps/backend/src/services/fulfillment/executors/kuaishou-cloud-executor.ts index 9551e4ae..9260f921 100644 --- a/apps/backend/src/services/fulfillment/executors/kuaishou-cloud-executor.ts +++ b/apps/backend/src/services/fulfillment/executors/kuaishou-cloud-executor.ts @@ -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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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, + } +} diff --git a/apps/backend/src/services/fulfillment/executors/registry.test.ts b/apps/backend/src/services/fulfillment/executors/registry.test.ts new file mode 100644 index 00000000..dcace942 --- /dev/null +++ b/apps/backend/src/services/fulfillment/executors/registry.test.ts @@ -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) +}) diff --git a/apps/backend/src/services/fulfillment/executors/registry.ts b/apps/backend/src/services/fulfillment/executors/registry.ts index 697754ac..3b00ab3f 100644 --- a/apps/backend/src/services/fulfillment/executors/registry.ts +++ b/apps/backend/src/services/fulfillment/executors/registry.ts @@ -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) +} diff --git a/apps/backend/src/services/fulfillment/executors/types.ts b/apps/backend/src/services/fulfillment/executors/types.ts index ad4e62bb..a64e4a88 100644 --- a/apps/backend/src/services/fulfillment/executors/types.ts +++ b/apps/backend/src/services/fulfillment/executors/types.ts @@ -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 resolveDeliveryLink?: (task: TaskRow) => Promise + /** lewan:准备绑定资源(虚拟号 / bindUrl) */ + prepareBinding?: ( + task: TaskRow, + options?: FulfillmentActionOptions, + ) => Promise + /** lewan:换绑角色 */ + rebindRole?: ( + task: TaskRow, + options?: FulfillmentActionOptions, + ) => Promise + /** lewan:刷新角色信息 */ + refreshRole?: ( + task: TaskRow, + options?: FulfillmentActionOptions, + ) => Promise + /** lewan:确认角色(UID 匹配后落 ROLE_CONFIRMED) */ + confirmRole?: ( + task: TaskRow, + options?: FulfillmentActionOptions, + ) => Promise + /** lewan:一键兑换(状态机 + dispatch) */ + redeemTask?: ( + task: TaskRow, + options?: FulfillmentActionOptions, + ) => Promise + /** lewan:后台/系统发货(不经 claim 一键兑换) */ + dispatchTask?: ( + task: TaskRow, + options?: FulfillmentActionOptions, + ) => Promise + /** + * lewan:退还虚拟号。 + * 默认可顺带核销;admin 清理请传 consumeIndustryVoucher: false。 + */ + returnNumber?: ( + task: TaskRow, + options?: FulfillmentActionOptions, + ) => Promise } export function normalizeExecutorKey(value: unknown): FulfillmentExecutorKey { diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/confirm-role.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/confirm-role.ts new file mode 100644 index 00000000..de9677ab --- /dev/null +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/confirm-role.ts @@ -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 { + 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 } +} diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/dispatch-failure.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/dispatch-failure.ts new file mode 100644 index 00000000..5bbcb159 --- /dev/null +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/dispatch-failure.ts @@ -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) + : 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) +} diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/dispatch-fulfillment.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/dispatch-fulfillment.ts new file mode 100644 index 00000000..4d902e15 --- /dev/null +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/dispatch-fulfillment.ts @@ -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) +} diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/dispatch-stock.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/dispatch-stock.ts new file mode 100644 index 00000000..c30d6a2f --- /dev/null +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/dispatch-stock.ts @@ -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>() + +/** + * 计算库存缺口 → 必要时采购 → 按数量逐次 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> + 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( + keys: string[], + callback: () => Promise, +): Promise { + const normalizedKeys = [ + ...new Set(keys.map((key) => String(key || '').trim()).filter(Boolean)), + ].sort() + + async function run(index: number): Promise { + 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(key: string, callback: () => Promise): Promise { + const previous = cloudSkuDispatchLocks.get(key) || Promise.resolve() + let release: () => void = () => {} + const current = new Promise((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) +} diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/dispatch-types.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/dispatch-types.ts new file mode 100644 index 00000000..0e7efc7a --- /dev/null +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/dispatch-types.ts @@ -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, + }, + ] +} diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/index.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/index.ts index 80f2407f..35db4be4 100644 --- a/apps/backend/src/services/fulfillment/kuaishou-cloud/index.ts +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/index.ts @@ -1,6 +1,9 @@ /** - * kuaishou-lewan(cloud)履约入口。 - * 实现已按职责拆分到 prepare / rebind / probe / refresh-role 等模块。 + * kuaishou-lewan(业务名)履约入口。 + * 模块路径仍为 kuaishou-cloud;executor_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' diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/mock-helpers.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/mock-helpers.ts new file mode 100644 index 00000000..5c02d681 --- /dev/null +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/mock-helpers.ts @@ -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 | 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, + ) +} diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/redeem-fulfillment.test.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/redeem-fulfillment.test.ts new file mode 100644 index 00000000..0297df59 --- /dev/null +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/redeem-fulfillment.test.ts @@ -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' }, + }, + ), + /不一致/, + ) +}) diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/redeem-fulfillment.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/redeem-fulfillment.ts new file mode 100644 index 00000000..da40a0ca --- /dev/null +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/redeem-fulfillment.ts @@ -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 { + 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 + } +} diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/return-fulfillment.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/return-fulfillment.ts new file mode 100644 index 00000000..c0dafc61 --- /dev/null +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/return-fulfillment.ts @@ -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>, 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) +} diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud/task-finalization.ts b/apps/backend/src/services/fulfillment/kuaishou-cloud/task-finalization.ts index 02a0d56f..9474a4a0 100644 --- a/apps/backend/src/services/fulfillment/kuaishou-cloud/task-finalization.ts +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud/task-finalization.ts @@ -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>(); - -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) - : 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>; - 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(keys: string[], callback: () => Promise): Promise { - const normalizedKeys = [...new Set(keys.map((key) => String(key || "").trim()).filter(Boolean))].sort(); - - async function run(index: number): Promise { - const key = normalizedKeys[index]; - if (!key) { - return callback(); - } - - return withCloudSkuDispatchLock(key, () => run(index + 1)); - } - - return run(0); -} - -async function withCloudSkuDispatchLock(key: string, callback: () => Promise): Promise { - const previous = cloudSkuDispatchLocks.get(key) || Promise.resolve(); - let release: () => void = () => {}; - const current = new Promise((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' diff --git a/docs/履约配置/工程债-lewan拆分与短链.md b/docs/履约配置/工程债-lewan拆分与短链.md index 5d9313d0..f64b003b 100644 --- a/docs/履约配置/工程债-lewan拆分与短链.md +++ b/docs/履约配置/工程债-lewan拆分与短链.md @@ -1,5 +1,14 @@ # 工程债:lewan 拆分、默认角色、短链 +## 命名约定 + +| 叫法 | 用途 | +| --- | --- | +| lewan / kuaishou-lewan | 业务/文案 | +| `kuaishou_ct_assisted` | `executor_key` | +| `kuaishou-cloud/` | 履约领域模块路径 | +| `cloudtentacles` | 外部发货平台 API | + ## lewan 模块拆分(已做) ```text @@ -13,11 +22,23 @@ kuaishou-cloud/ delivery-plan.ts ensure-claim-link.ts dispatch-role-sync.ts - task-finalization.ts # dispatch / return(仍厚,后续可再拆) + redeem-fulfillment.ts # 一键兑换用例 + confirm-role.ts # 确认角色用例 + dispatch-fulfillment.ts # 发货编排 + dispatch-stock.ts # 库存/采购/下发 + dispatch-failure.ts # 发货失败落库 + return-fulfillment.ts # 退号 + 电子凭证核销收尾 + mock-helpers.ts + task-finalization.ts # 兼容 re-export domain.ts + +executors/ + kuaishou-cloud-executor # preparePaid + prepareBinding + rebind/refresh/confirm + redeem + dispatch + return + registry # prepare/rebind/refresh/confirm/redeem/dispatch/return 统一入口 ``` -对外 import 路径仍为 `.../kuaishou-cloud/index.js`。 +对外 import 路径仍为 `.../kuaishou-cloud/index.js`。 +claim / admin 应优先走 `executors/registry`,不要直连子模块动作。 ## 默认角色