diff --git a/apps/backend/src/services/admin/admin-write-service.js b/apps/backend/src/services/admin/admin-write-service.js index 019cc744..9deebe15 100644 --- a/apps/backend/src/services/admin/admin-write-service.js +++ b/apps/backend/src/services/admin/admin-write-service.js @@ -1,1710 +1,30 @@ -// @ts-check - -import { - createInventoryItems, - invalidateInventoryItem, - listInventoryItems, - markInventoryItemDelivered, - releaseReservedInventoryItem, -} from '../../repositories/inventory-repo.js' -import { updateClaimToken } from '../../repositories/claim-token-repo.js' -import { listOrderItemsByOrderId } from '../../repositories/order-item-repo.js' -import { getOrderById } from '../../repositories/order-repo.js' -import { updateTask } from '../../repositories/task-repo.js' -import { - getTaskInventoryBindingById, - listTaskInventoryBindingsByTaskId, -} from '../../repositories/task-inventory-binding-repo.js' -import { createTaskEvent } from '../../repositories/task-event-repo.js' -import { getWebhookEventById } from '../../repositories/webhook-event-repo.js' -import { buildClaimUrl, createTaskClaimToken } from '../claim/claim-service.js' -import { confirmClaimRoleForAdminTask, redeemClaimTaskForAdminTask } from '../claim/claim-session-service.js' -import { reserveInventoryForTask } from '../order/inventory-service.js' -import { normalizeProductName } from '../order/product-match-service.js' -import { replayAgisoTradeWebhookEvent } from '../order/webhook-service.js' -import { ensureAgisoXianyuAutoDeliveryForDeliveredTask } from '../platforms/agiso/xianyu/auto-delivery-service.js' -import { getCloudtentaclesSourceConfig } from '../platforms/cloudtentacles/source-config-service.js' -import { getCloudtentaclesSessionState } from '../platforms/cloudtentacles/session-state-service.js' -import { - buyCloudtentaclesSku, - getCloudtentaclesAsset, - listCloudtentaclesSku, - useCloudtentaclesSku, -} from '../platforms/cloudtentacles/catalog-service.js' -import { getCloudtentaclesKnapsack } from '../platforms/cloudtentacles/knapsack-service.js' -import { consumeKuaishouEticket } from '../platforms/kuaishou-eticket/consume-service.js' -import { - getKuaishouEticketSourceConfig, - resolveKuaishouEticketShopConfig, -} from '../platforms/kuaishou-eticket/source-config-service.js' -import { - appointCloudtentaclesVirtualNumber, - backCloudtentaclesVirtualNumber, - getCloudtentaclesBindInfo, - fetchCloudtentaclesVirtualNumberCode, - generateCloudtentaclesLoginCode, - getCloudtentaclesBindUrl, - verifyCloudtentaclesLoginCode, -} from '../platforms/cloudtentacles/virtual-number-service.js' -import { closeTencentBrowserSession } from '../session/session.js' -import { createHttpError } from '../../utils/http.js' -import { nowIso } from '../../utils/time.js' -import { - canViewerConfirmAssistedRole, - canViewerRedeemAssistedTask, - createAdminViewerContext, - getTaskPrimaryClaimTokenId, - getTaskPrimaryInventoryItemId, - isAssistedClaimTask, - isManualDispatchTask, - parseTaskContext, -} from './admin-read-shared-helpers.js' -import { - getRequiredInventoryItem, - mapAdminInventoryListItem, -} from './admin-inventory-read-helpers.js' -import { - getRequiredTask, - mapTaskActionPayload, -} from './admin-task-read-helpers.js' - -/** @typedef {import('../../types/admin-read-inputs.js').AdminEntityIdInput} AdminEntityIdInput */ -/** @typedef {import('../../types/admin-read-inputs.js').AdminViewerSessionInput} AdminViewerSessionInput */ -/** @typedef {import('../../types/admin-write-inputs.js').AdminInventoryCreateInput} AdminInventoryCreateInput */ -/** @typedef {import('../../types/admin-write-inputs.js').AdminInventoryImportInput} AdminInventoryImportInput */ -/** @typedef {import('../../types/admin-write-inputs.js').AdminInventoryImportRowInput} AdminInventoryImportRowInput */ -/** @typedef {import('../../types/admin-write-inputs.js').AdminInventoryInvalidateInput} AdminInventoryInvalidateInput */ -/** @typedef {import('../../types/admin-write-models.js').AdminInventoryImportResponse} AdminInventoryImportResponse */ -/** @typedef {import('../../types/admin-write-models.js').AdminInventoryMutationResponse} AdminInventoryMutationResponse */ -/** @typedef {import('../../types/admin-write-models.js').AdminTaskActionResponse} AdminTaskActionResponse */ -/** @typedef {import('../../types/admin-write-models.js').AdminTaskBindingReleaseResponse} AdminTaskBindingReleaseResponse */ -/** @typedef {import('../../types/admin-write-models.js').AdminTaskManualDispatchResponse} AdminTaskManualDispatchResponse */ -/** @typedef {import('../../types/admin-write-models.js').AdminWebhookReplayResponse} AdminWebhookReplayResponse */ -/** @typedef {import('../../types/admin-write-inputs.js').AdminTaskKuaishouCloudDispatchInput} AdminTaskKuaishouCloudDispatchInput */ - -const KUAISHOU_CLOUD_FIXED_VN_KEY = '1' - -/** @returns {Promise} */ -/** @param {AdminInventoryCreateInput} [payload] */ -export async function createAdminInventoryItem(payload = /** @type {AdminInventoryCreateInput} */ ({})) { - const skuCode = String(payload.skuCode || '').trim() - const displayValue = String(payload.displayValue || '').trim() - const batchNo = String(payload.batchNo || '').trim() - const credentialType = String(payload.credentialType || 'tencent_code').trim() || 'tencent_code' - const inventoryGroupCode = String(payload.inventoryGroupCode || '').trim() - - if (!skuCode || !displayValue) { - throw createHttpError('缺少 skuCode 或 displayValue', { - statusCode: 400, - errorCode: 'admin_inventory_create_invalid', - }) - } - - const now = nowIso() - const created = await createInventoryItems([ - { - skuCode, - displayValue, - batchNo, - credentialType, - inventoryGroupCode, - createdAt: now, - updatedAt: now, - }, - ]) - - if (created === 0) { - throw createHttpError('库存凭据已存在,不能重复新增', { - statusCode: 409, - errorCode: 'admin_inventory_duplicate', - }) - } - - const { items } = await listInventoryItems({ - page: 1, - pageSize: 20, - skuCode, - inventoryGroupCode, - }) - const createdItem = items.find((item) => item.display_value === displayValue && item.inventory_group_code === inventoryGroupCode) || null - - return { - inventoryItem: createdItem ? await mapAdminInventoryListItem(createdItem) : null, - } -} - -/** @returns {Promise} */ -/** @param {AdminInventoryImportInput} [payload] */ -export async function importAdminInventoryItems(payload = /** @type {AdminInventoryImportInput} */ ({})) { - const rows = normalizeInventoryImportRows(payload) - - if (rows.length === 0) { - throw createHttpError('没有可导入的库存凭据数据', { - statusCode: 400, - errorCode: 'admin_inventory_import_empty', - }) - } - - const now = nowIso() - const normalizedRows = rows.map((row) => ({ - skuCode: row.skuCode, - batchNo: row.batchNo, - displayValue: row.displayValue, - credentialType: row.credentialType, - inventoryGroupCode: row.inventoryGroupCode, - createdAt: now, - updatedAt: now, - })) - - const created = await createInventoryItems(normalizedRows) - - return { - total: normalizedRows.length, - created, - duplicated: normalizedRows.length - created, - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} inventoryItemId */ -export async function releaseAdminInventoryItem(inventoryItemId) { - const inventoryItem = await getRequiredInventoryItem(inventoryItemId) - - if (inventoryItem.status !== 'reserved') { - throw createHttpError('当前库存项不是预占状态,不能释放', { - statusCode: 409, - errorCode: 'admin_inventory_release_not_allowed', - }) - } - - const updated = await releaseReservedInventoryItem(inventoryItem.id, nowIso()) - - return { - inventoryItem: await mapAdminInventoryListItem(updated), - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} inventoryItemId */ -/** @param {AdminInventoryInvalidateInput} [payload] */ -export async function invalidateAdminInventoryItem( - inventoryItemId, - payload = /** @type {AdminInventoryInvalidateInput} */ ({}), -) { - const inventoryItem = await getRequiredInventoryItem(inventoryItemId) - - if (inventoryItem.status !== 'available') { - throw createHttpError('只有可用库存项才能作废,请先释放预占', { - statusCode: 409, - errorCode: 'admin_inventory_invalidate_not_allowed', - }) - } - - const reason = String(payload.reason || '').trim() || '后台手动作废' - const updated = await invalidateInventoryItem(inventoryItem.id, reason, nowIso()) - - return { - inventoryItem: await mapAdminInventoryListItem(updated), - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} eventId */ -export async function replayAdminWebhookEvent(eventId) { - const event = await getWebhookEventById(Number(eventId)) - - if (!event) { - throw createHttpError('Webhook 事件不存在', { - statusCode: 404, - errorCode: 'admin_webhook_not_found', - }) - } - - if (String(event.provider || event.platform || '').trim() !== 'agiso') { - throw createHttpError('当前只支持重放 agiso webhook', { - statusCode: 409, - errorCode: 'admin_webhook_replay_not_supported', - }) - } - - const result = await replayAgisoTradeWebhookEvent(event) - - return { - eventId: event.id, - replayed: true, - result, - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -export async function releaseAdminTaskInventory(taskId) { - const task = await getRequiredTask(taskId) - const primaryInventoryItemId = getTaskPrimaryInventoryItemId(task) - - if (!primaryInventoryItemId) { - throw createHttpError('当前任务没有预占库存项', { - statusCode: 409, - errorCode: 'admin_task_no_reserved_inventory', - }) - } - - if (task.task_status === 'redeemed') { - throw createHttpError('已兑换任务不能释放库存项', { - statusCode: 409, - errorCode: 'admin_task_release_not_allowed', - }) - } - - const now = nowIso() - await releaseReservedInventoryItem(primaryInventoryItemId, now) - const updatedTask = await updateTask(task.id, { - task_status: 'waiting_inventory', - inventory_status: 'pending', - last_error: '已手动释放预占库存项', - updated_at: now, - }) - - return { - task: mapTaskActionPayload(updatedTask), - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -/** @param {AdminEntityIdInput} bindingId */ -export async function releaseAdminTaskInventoryBinding(taskId, bindingId) { - const task = await getRequiredTask(taskId) - const binding = await getTaskInventoryBindingById(Number(bindingId)) - - if (!binding || Number(binding.task_id) !== Number(task.id)) { - throw createHttpError('任务库存绑定不存在', { - statusCode: 404, - errorCode: 'admin_task_inventory_binding_not_found', - }) - } - - if (String(binding.binding_status || '').trim() !== 'reserved') { - throw createHttpError('当前库存绑定不是预占状态,不能释放', { - statusCode: 409, - errorCode: 'admin_task_inventory_binding_release_not_allowed', - }) - } - - if (task.task_status === 'redeemed') { - throw createHttpError('已兑换任务不能释放库存绑定', { - statusCode: 409, - errorCode: 'admin_task_release_not_allowed', - }) - } - - const now = nowIso() - await releaseReservedInventoryItem(binding.inventory_item_id, now) - - const remainingBindings = await listTaskInventoryBindingsByTaskId(task.id) - const activeBindings = remainingBindings.filter((item) => ['reserved', 'consumed'].includes(String(item.binding_status || '').trim())) - const hasReservedBindings = activeBindings.some((item) => String(item.binding_status || '').trim() === 'reserved') - const hasConsumedBindings = activeBindings.some((item) => String(item.binding_status || '').trim() === 'consumed') - const nextInventoryStatus = hasReservedBindings ? 'reserved' : (hasConsumedBindings ? 'consumed' : 'pending') - const nextTaskStatus = !activeBindings.length && !['closed', 'expired'].includes(String(task.task_status || '').trim()) - ? 'waiting_inventory' - : task.task_status - - if (!hasReservedBindings) { - const primaryClaimTokenId = getTaskPrimaryClaimTokenId(task) - if (primaryClaimTokenId) { - await updateClaimToken(primaryClaimTokenId, { - status: 'revoked', - updated_at: now, - }) - } - } - - const updatedTask = await updateTask(task.id, { - task_status: nextTaskStatus, - inventory_status: nextInventoryStatus, - last_error: !activeBindings.length ? '已手动释放预占库存绑定' : (task.last_error || ''), - updated_at: now, - }) - - await createTaskEvent(task.id, 'inventory_binding_released', { - bindingId: Number(binding.id), - inventoryItemId: Number(binding.inventory_item_id), - roleKey: String(binding.role_key || '').trim(), - }, now) - - return { - task: mapTaskActionPayload(updatedTask), - bindingId: Number(binding.id), - inventoryItemId: Number(binding.inventory_item_id), - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -/** @param {AdminViewerSessionInput | null} [session] */ -export async function regenerateAdminTaskClaimLink(taskId, session = null) { - const task = await getRequiredTask(taskId) - const primaryClaimTokenId = getTaskPrimaryClaimTokenId(task) - const viewerContext = createAdminViewerContext(session) - - if (isManualDispatchTask(task)) { - throw createHttpError('人工履约任务不需要领取链接,请直接回写人工履约结果', { - statusCode: 409, - errorCode: 'admin_task_manual_dispatch_claim_not_allowed', - }) - } - - if (['redeemed', 'closed'].includes(task.task_status)) { - throw createHttpError('当前任务状态不允许重新生成领取链接', { - statusCode: 409, - errorCode: 'admin_task_regenerate_not_allowed', - }) - } - - if (!viewerContext.canManageTaskLifecycle && !isAssistedClaimTask(task)) { - throw createHttpError('当前账号只能重发半自动客服任务的领取链接', { - statusCode: 403, - errorCode: 'admin_task_regenerate_permission_denied', - }) - } - - const now = nowIso() - if (primaryClaimTokenId) { - await updateClaimToken(primaryClaimTokenId, { - status: 'revoked', - updated_at: now, - }) - } - - const claimToken = await createTaskClaimToken(task.id) - const updatedTask = await updateTask(task.id, { - claim_token: claimToken.token, - claim_expires_at: claimToken.expired_at, - task_status: 'link_generated', - last_error: '', - updated_at: now, - }) - - return { - task: mapTaskActionPayload(updatedTask), - claimUrl: claimToken.claimUrl, - token: claimToken.token, - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -/** @param {AdminViewerSessionInput | null} [session] */ -export async function confirmAdminTaskAssistedRole(taskId, session = null) { - const task = await getRequiredTask(taskId) - const viewerContext = createAdminViewerContext(session) - ensureViewerCanOperateAssistedTask(task, viewerContext, 'confirm') - await confirmClaimRoleForAdminTask(task.id) - const updatedTask = await getRequiredTask(task.id) - - return { - task: mapTaskActionPayload(updatedTask), - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -/** @param {AdminViewerSessionInput | null} [session] */ -export async function redeemAdminTaskAssisted(taskId, session = null) { - const task = await getRequiredTask(taskId) - const viewerContext = createAdminViewerContext(session) - ensureViewerCanOperateAssistedTask(task, viewerContext, 'redeem') - await redeemClaimTaskForAdminTask(task.id) - const updatedTask = await getRequiredTask(task.id) - - return { - task: mapTaskActionPayload(updatedTask), - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -export async function closeAdminTask(taskId) { - return closeAdminTaskWithDeps(taskId) -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -export async function closeAdminTaskWithDeps( - taskId, - { - getRequiredTask: getTask = getRequiredTask, - listTaskInventoryBindingsByTaskId: listBindings = listTaskInventoryBindingsByTaskId, - releaseReservedInventoryItem: releaseReserved = releaseReservedInventoryItem, - updateClaimToken: updateToken = updateClaimToken, - updateTask: updateTaskRecord = updateTask, - createTaskEvent: createEvent = createTaskEvent, - closeTencentBrowserSession: closeSession = closeTencentBrowserSession, - nowIso: getNowIso = nowIso, - } = {}, -) { - const task = await getTask(taskId) - - if (task.task_status === 'redeemed') { - throw createHttpError('已兑换任务不能关闭', { - statusCode: 409, - errorCode: 'admin_task_close_not_allowed', - }) - } - - const now = getNowIso() - const bindings = await listBindings(task.id) - const reservedBindings = bindings.filter((binding) => String(binding.binding_status || '').trim() === 'reserved') - const releasedInventoryItemIds = Array.from(new Set(reservedBindings.map((binding) => Number(binding.inventory_item_id)).filter((id) => id > 0))) - const hasConsumedBindings = bindings.some((binding) => String(binding.binding_status || '').trim() === 'consumed') - const primaryClaimTokenId = getTaskPrimaryClaimTokenId(task) - let browserSessionClosed = false - - if (primaryClaimTokenId) { - await updateToken(primaryClaimTokenId, { - status: 'revoked', - expired_at: now, - updated_at: now, - }) - } - - if (task.browser_session_id) { - try { - await closeSession(task.browser_session_id) - browserSessionClosed = true - } catch (error) { - if (!isRecoverableTaskSessionCloseError(error)) { - throw error - } - } - } - - for (const inventoryItemId of releasedInventoryItemIds) { - await releaseReserved(inventoryItemId, now) - } - - const closeReasonParts = ['已手动关闭任务'] - if (primaryClaimTokenId) { - closeReasonParts.push('领取链接已失效') - } - if (releasedInventoryItemIds.length > 0) { - closeReasonParts.push('预占库存已释放') - } - - const updatedTask = await updateTaskRecord(task.id, { - task_status: 'closed', - inventory_status: hasConsumedBindings ? 'consumed' : 'pending', - delivery_status: 'closed', - user_action_status: 'closed', - claim_expires_at: primaryClaimTokenId ? now : task.claim_expires_at, - browser_session_id: '', - last_error: task.last_error || closeReasonParts.join(','), - updated_at: now, - }) - - await createEvent(task.id, 'task_closed', { - claimTokenRevoked: Boolean(primaryClaimTokenId), - releasedInventoryItemIds, - releasedInventoryCount: releasedInventoryItemIds.length, - browserSessionClosed, - }, now) - - return { - task: mapTaskActionPayload(updatedTask), - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -export async function markAdminTaskManualReview(taskId) { - const task = await getRequiredTask(taskId) - const updatedTask = await updateTask(task.id, { - task_status: 'manual_review', - delivery_status: task.delivery_status || 'pending', - last_error: task.last_error || '已转人工处理', - updated_at: nowIso(), - }) - - return { - task: mapTaskActionPayload(updatedTask), - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -export async function retryAdminTask(taskId) { - const task = await getRequiredTask(taskId) - const now = nowIso() - - if (isManualDispatchTask(task)) { - throw createHttpError('人工履约任务不能走自动重试,请在详情页直接回写人工履约结果', { - statusCode: 409, - errorCode: 'admin_task_manual_dispatch_retry_not_allowed', - }) - } - - if (!['retry_pending', 'manual_review', 'waiting_inventory'].includes(task.task_status)) { - throw createHttpError('当前任务状态不允许重试', { - statusCode: 409, - errorCode: 'admin_task_retry_not_allowed', - }) - } - - const taskContext = parseTaskContext(task) - const primaryRequirement = taskContext.primaryRequirement || null - const orderItems = await listOrderItemsByOrderId(task.order_id) - const orderItem = orderItems.find((item) => item.id === task.order_item_id) || null - let reservedInventoryItemId = getTaskPrimaryInventoryItemId(task) - let claimTokenId = getTaskPrimaryClaimTokenId(task) - let nextStatus = 'link_generated' - let lastError = '' - let claimExpiresAt = getTaskClaimExpiresAt(task) - let claimUrl = '' - let token = '' - - if (!reservedInventoryItemId) { - const reserved = await reserveInventoryForTask({ - skuCode: orderItem?.sku_code || '', - taskId: task.id, - credentialType: primaryRequirement?.credentialType || 'tencent_code', - roleKey: primaryRequirement?.roleKey || 'primary_code', - inventoryGroupCodes: resolveTaskInventoryGroupCodes(task), - }) - - if (!reserved) { - nextStatus = 'waiting_inventory' - lastError = '库存不足,等待可用库存凭据' - } else { - reservedInventoryItemId = reserved.id - } - } - - if (nextStatus === 'link_generated' && !claimTokenId) { - const claimToken = await createTaskClaimToken(task.id) - claimTokenId = claimToken.id - claimExpiresAt = claimToken.expired_at - claimUrl = claimToken.claimUrl - token = claimToken.token - } - - const updatedTask = await updateTask(task.id, { - task_status: nextStatus, - inventory_status: reservedInventoryItemId ? 'reserved' : 'pending', - claim_token: token || task.claim_token || '', - claim_expires_at: claimExpiresAt, - last_error: lastError, - updated_at: now, - }) - - const response = { - task: mapTaskActionPayload(updatedTask), - } - - if (claimUrl) { - response.claimUrl = claimUrl - response.token = token - } - - return response -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -/** @param {{ outcome?: string, resultMessage?: string, deliveryReference?: string, deliveredCredential?: string }} [payload] */ -/** @param {AdminViewerSessionInput | null} [session] */ -export async function completeAdminTaskManualDispatch(taskId, payload = {}, session = null) { - const task = await getRequiredTask(taskId) - - if (!isManualDispatchTask(task)) { - throw createHttpError('当前任务不是人工履约任务', { - statusCode: 409, - errorCode: 'admin_task_not_manual_dispatch', - }) - } - - if (['redeemed', 'closed'].includes(task.task_status)) { - throw createHttpError('当前任务已经完结,不能重复回写人工履约结果', { - statusCode: 409, - errorCode: 'admin_task_manual_dispatch_already_completed', - }) - } - - const now = nowIso() - const outcome = normalizeManualDispatchOutcome(payload.outcome) - const resultMessage = String(payload.resultMessage || '').trim() - const deliveryReference = String(payload.deliveryReference || '').trim() - const deliveredCredential = String(payload.deliveredCredential || '').trim() - const context = parseTaskContext(task) - const inventoryItemId = getTaskPrimaryInventoryItemId(task) - const resultCode = outcome === 'failed' ? 'manual_dispatch_failed' : 'manual_dispatch_delivered' - const fallbackMessage = outcome === 'failed' ? '人工履约失败' : '人工履约已完成' - const nextTaskStatus = outcome === 'failed' ? 'closed' : 'redeemed' - const nextDeliveryStatus = outcome === 'failed' ? 'failed' : 'delivered' - const nextInventoryStatus = outcome === 'delivered' && inventoryItemId - ? 'consumed' - : task.inventory_status || 'not_required' - const manualDispatch = { - outcome, - deliveryReference, - deliveredCredential, - resultMessage: resultMessage || fallbackMessage, - completedAt: now, - completedBy: session - ? { - userId: Number(session.userId || 0), - username: String(session.username || ''), - role: String(session.role || ''), - } - : null, - } - - if (outcome === 'delivered' && inventoryItemId) { - await markInventoryItemDelivered(inventoryItemId, now) - } - - const updatedTask = await updateTask(task.id, { - task_status: nextTaskStatus, - inventory_status: nextInventoryStatus, - delivery_status: nextDeliveryStatus, - result_code: resultCode, - result_message: resultMessage || fallbackMessage, - user_action_status: 'not_required', - last_error: outcome === 'failed' ? (resultMessage || fallbackMessage) : '', - redeemed_at: outcome === 'delivered' ? now : task.redeemed_at || null, - context_json: JSON.stringify({ - ...context, - manualDispatch, - }), - updated_at: now, - }) - - await createTaskEvent(task.id, 'manual_dispatch_completed', { - outcome, - resultCode, - resultMessage: resultMessage || fallbackMessage, - deliveryReference, - deliveredCredentialMasked: maskCode(deliveredCredential), - completedBy: manualDispatch.completedBy, - }, now) - - const taskAfterAutoDelivery = outcome === 'delivered' - ? (await ensureAgisoXianyuAutoDeliveryForDeliveredTask({ - order: await getOrderById(task.order_id), - task: updatedTask, - trigger: 'manual_dispatch_completed', - })).task || updatedTask - : updatedTask - - return { - outcome, - task: mapTaskActionPayload(taskAfterAutoDelivery), - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -/** @param {AdminViewerSessionInput | null} [session] */ -export async function prepareAdminTaskKuaishouCloudFulfillment(taskId, session = null) { - const task = await getRequiredTask(taskId) - const now = nowIso() - const viewerContext = createAdminViewerContext(session) - - if (!viewerContext.canManageTaskLifecycle) { - throw createHttpError('当前账号没有权限准备绑定资源', { - statusCode: 403, - errorCode: 'admin_task_prepare_kuaishou_cloud_forbidden', - }) - } - - if (!isKuaishouCloudTask(task)) { - throw createHttpError('当前任务不是快手 cloud 履约任务', { - statusCode: 409, - errorCode: 'admin_task_not_kuaishou_cloud', - }) - } - - const taskContext = parseTaskContext(task) - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) - const cloudContext = resolvePersistedCloudtentaclesContext() - const [knapsack, skuList] = await Promise.all([ - getCloudtentaclesKnapsack(cloudContext), - listCloudtentaclesSku(cloudContext), - ]) - const resolvedBinding = resolveKuaishouCloudBindingResources(flow, { - skuItems: Array.isArray(skuList.items) ? skuList.items : [], - knapsackItems: Array.isArray(knapsack.items) ? knapsack.items : [], - }) - const vnKeyCandidates = resolveKuaishouCloudVnKeyCandidates({ - flow, - binding: resolvedBinding, - }) - - if (!resolvedBinding.skuId || vnKeyCandidates.length === 0) { - throw createHttpError('当前任务缺少可用 cloud 资源,且无法根据商品名自动解析 SKU / VN Key', { - statusCode: 409, - errorCode: 'admin_task_kuaishou_cloud_missing_binding_config', - }) - } - - const flowWithResolvedBinding = { - ...flow, - binding: { - ...flow.binding, - skuId: resolvedBinding.skuId, - skuName: resolvedBinding.skuName, - vnKey: KUAISHOU_CLOUD_FIXED_VN_KEY, - }, - } - const knapsackItem = resolvedBinding.knapsackItem - let usedKnapsack = Number(knapsackItem?.count || 0) > 0 - let purchaseTriggered = false - let assetBefore = 0 - let assetAfter = 0 - const claimLinkState = await ensureTaskClaimLink(task) - - if (!usedKnapsack) { - if (!flowWithResolvedBinding.purchase.autoBuyEnabled) { - throw createHttpError('背包中没有现成库存,且当前配置未开启自动购买', { - statusCode: 409, - errorCode: 'admin_task_kuaishou_cloud_auto_buy_disabled', - }) - } - - const asset = await getCloudtentaclesAsset(cloudContext) - const targetSku = resolvedBinding.skuItem - - if (!targetSku) { - throw createHttpError(`cloudtentacles 未找到 SKU ${flowWithResolvedBinding.binding.skuId}`, { - statusCode: 404, - errorCode: 'admin_task_kuaishou_cloud_sku_not_found', - }) - } - - assetBefore = Number(asset.asset || 0) || 0 - const targetPrice = Number(targetSku.price || 0) || 0 - const requiredAsset = targetPrice + flowWithResolvedBinding.purchase.minAssetReserve - if (assetBefore < requiredAsset) { - throw createHttpError(`余额不足,当前 ${assetBefore},至少需要 ${requiredAsset}`, { - statusCode: 409, - errorCode: 'admin_task_kuaishou_cloud_asset_not_enough', - }) - } - - await buyCloudtentaclesSku({ - ...cloudContext, - id: flowWithResolvedBinding.binding.skuId, - count: 1, - }) - purchaseTriggered = true - - const assetResult = await getCloudtentaclesAsset(cloudContext) - assetAfter = Number(assetResult.asset || 0) || 0 - } - - const preparedBinding = await prepareKuaishouCloudBindResourceWithFallback({ - cloudContext, - vnKeyCandidates, - }) - const vnId = preparedBinding.vnId - const vnPhone = preparedBinding.vnPhone - const bindUrlResult = { - bindUrl: preparedBinding.bindUrl, - } - - const nextContext = { - ...taskContext, - kuaishouCloudFulfillment: { - ...flowWithResolvedBinding, - binding: { - ...flowWithResolvedBinding.binding, - vnKey: preparedBinding.vnKey, - prepareStatus: 'ready', - vnId, - vnPhone, - bindUrl: bindUrlResult.bindUrl, - bindPreparedAt: now, - roleName: '', - roleId: '', - }, - role: { - status: 'pending', - name: '', - rid: '', - refreshedAt: null, - errorMessage: '', - rawInfo: null, - }, - purchase: { - ...flowWithResolvedBinding.purchase, - usedKnapsack, - purchaseTriggered, - assetBefore, - assetAfter, - purchaseAt: purchaseTriggered ? now : flowWithResolvedBinding.purchase.purchaseAt, - }, - }, - } - - const updatedTask = await updateTask(task.id, { - task_status: 'waiting_binding', - inventory_status: 'not_required', - user_action_status: 'pending', - claim_token: claimLinkState.token || task.claim_token || '', - claim_expires_at: claimLinkState.expiredAt || getTaskClaimExpiresAt(task), - last_error: '', - context_json: JSON.stringify(nextContext), - updated_at: now, - }) - - await createTaskEvent(task.id, 'kuaishou_cloud_binding_prepared', { - skuId: flowWithResolvedBinding.binding.skuId, - skuName: flowWithResolvedBinding.binding.skuName, - vnKey: preparedBinding.vnKey, - vnId, - vnPhoneMasked: maskPhone(vnPhone), - bindUrl: bindUrlResult.bindUrl, - purchaseTriggered, - usedKnapsack, - resolvedByName: resolvedBinding.resolvedByName, - }, now) - - return { - task: mapTaskActionPayload(updatedTask), - claimUrl: claimLinkState.claimUrl, - token: claimLinkState.token, - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -/** @param {AdminTaskKuaishouCloudDispatchInput} [payload] */ -/** @param {AdminViewerSessionInput | null} [session] */ -export async function dispatchAdminTaskKuaishouCloudFulfillment(taskId, payload = {}, session = null) { - const task = await getRequiredTask(taskId) - const now = nowIso() - const viewerContext = createAdminViewerContext(session) - - if (!viewerContext.canManageTaskLifecycle) { - throw createHttpError('当前账号没有权限执行发货', { - statusCode: 403, - errorCode: 'admin_task_dispatch_kuaishou_cloud_forbidden', - }) - } - - if (!isKuaishouCloudTask(task)) { - throw createHttpError('当前任务不是快手 cloud 履约任务', { - statusCode: 409, - errorCode: 'admin_task_not_kuaishou_cloud', - }) - } - - const taskContext = parseTaskContext(task) - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) - const cloudContext = resolvePersistedCloudtentaclesContext() - - if (!flow.binding.skuId || !flow.binding.vnId || !flow.binding.vnPhone) { - throw createHttpError('当前任务还没有准备好绑定资源,请先准备绑定资源', { - statusCode: 409, - errorCode: 'admin_task_kuaishou_cloud_not_prepared', - }) - } - - const ticketCode = String(payload.ticketCode || '').trim() - const persistedTicketCode = String(flow.ticket.code || '').trim() - - if (!persistedTicketCode && !ticketCode) { - throw createHttpError('客户还没有在领取页提交核销码,暂时不能直接发货', { - statusCode: 409, - errorCode: 'admin_task_kuaishou_cloud_missing_ticket_code', - }) - } - - const dispatchResult = await useCloudtentaclesSku({ - ...cloudContext, - id: flow.binding.skuId, - virtualNumberId: flow.binding.vnId, - phone: flow.binding.vnPhone, - }) - - const nextContext = { - ...taskContext, - kuaishouCloudFulfillment: { - ...flow, - ticket: { - ...flow.ticket, - code: ticketCode || persistedTicketCode, - capturedAt: ticketCode ? now : flow.ticket.capturedAt, - capturedBy: ticketCode && session - ? { - userId: Number(session.userId || 0) || 0, - username: String(session.username || '').trim(), - role: String(session.role || '').trim(), - } - : flow.ticket.capturedBy, - }, - dispatch: { - ...flow.dispatch, - status: 'success', - dispatchAt: now, - dispatchBy: session - ? { - userId: Number(session.userId || 0) || 0, - username: String(session.username || '').trim(), - role: String(session.role || '').trim(), - } - : null, - sendType: Number(dispatchResult.sendType || 0) || 0, - note: String(dispatchResult.note || dispatchResult.responseMessage || '').trim(), - }, - }, - } - - const updatedTask = await updateTask(task.id, { - task_status: 'dispatched_pending_return', - delivery_status: 'delivered', - result_code: 'kuaishou_cloud_dispatched', - result_message: String(dispatchResult.responseMessage || dispatchResult.note || 'cloudtentacles 发货成功').trim(), - last_error: '', - context_json: JSON.stringify(nextContext), - updated_at: now, - }) - - await createTaskEvent(task.id, 'kuaishou_cloud_dispatched', { - ticketCodeMasked: maskCode(ticketCode || persistedTicketCode), - skuId: flow.binding.skuId, - vnId: flow.binding.vnId, - vnPhoneMasked: maskPhone(flow.binding.vnPhone), - sendType: dispatchResult.sendType, - note: dispatchResult.note, - }, now) - - return { - task: mapTaskActionPayload(updatedTask), - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -/** @param {AdminViewerSessionInput | null} [session] */ -export async function refreshAdminTaskKuaishouCloudRoleInfo(taskId, session = null) { - const task = await getRequiredTask(taskId) - const now = nowIso() - const viewerContext = createAdminViewerContext(session) - - if (!viewerContext.canOperateAssistedTask) { - throw createHttpError('当前账号没有权限刷新角色信息', { - statusCode: 403, - errorCode: 'admin_task_refresh_kuaishou_cloud_role_forbidden', - }) - } - - if (!isKuaishouCloudTask(task)) { - throw createHttpError('当前任务不是快手 cloud 履约任务', { - statusCode: 409, - errorCode: 'admin_task_not_kuaishou_cloud', - }) - } - - const taskContext = parseTaskContext(task) - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) - if (!flow.binding.vnId || !flow.binding.vnKey) { - throw createHttpError('当前任务还没有可查询的绑定角色信息,请先准备绑定资源', { - statusCode: 409, - errorCode: 'admin_task_kuaishou_cloud_missing_bind_info_context', - }) - } - - const cloudContext = resolvePersistedCloudtentaclesContext() - 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), - } -} - -/** @returns {Promise} */ -/** @param {AdminEntityIdInput} taskId */ -/** @param {AdminViewerSessionInput | null} [session] */ -export async function returnNumberAdminTaskKuaishouCloudFulfillment(taskId, session = null) { - const task = await getRequiredTask(taskId) - const now = nowIso() - const viewerContext = createAdminViewerContext(session) - - if (!viewerContext.canManageTaskLifecycle) { - throw createHttpError('当前账号没有权限退还号码', { - statusCode: 403, - errorCode: 'admin_task_return_kuaishou_cloud_forbidden', - }) - } - - if (!isKuaishouCloudTask(task)) { - throw createHttpError('当前任务不是快手 cloud 履约任务', { - statusCode: 409, - errorCode: 'admin_task_not_kuaishou_cloud', - }) - } - - const taskContext = parseTaskContext(task) - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) - const cloudContext = resolvePersistedCloudtentaclesContext() - - if (!flow.binding.vnId || !flow.binding.vnKey) { - throw createHttpError('当前任务缺少可退还的虚拟号信息', { - statusCode: 409, - errorCode: 'admin_task_kuaishou_cloud_missing_return_context', - }) - } - - await backCloudtentaclesVirtualNumber({ - ...cloudContext, - key: flow.binding.vnKey, - id: flow.binding.vnId, - }) - - const order = await getOrderById(task.order_id) - const ticketCode = String(flow.ticket.code || '').trim() - const shopId = String(order?.shop_id || '').trim() - const eticketSource = getKuaishouEticketSourceConfig() - const shopConfig = resolveKuaishouEticketShopConfig({ - shopId, - shopName: String(order?.shop_name || '').trim(), - }) - let consumeStatus = 'pending' - let consumeErrorMessage = '' - let consumedAt = null - let nextTaskStatus = 'completed' - let nextResultCode = 'kuaishou_cloud_completed' - let nextResultMessage = 'cloudtentacles 发货、退号并完成快手核销' - - if (!order) { - consumeStatus = 'failed' - consumeErrorMessage = '任务关联订单不存在,无法执行快手核销' - } else if (!ticketCode) { - consumeStatus = 'failed' - consumeErrorMessage = '客户未提交有效核销码,无法执行快手核销' - } else if (!shopConfig || shopConfig.enabled === false || !String(shopConfig.cookie || '').trim()) { - consumeStatus = 'failed' - consumeErrorMessage = '订单对应快手小店缺少可用 Cookie,无法执行快手核销' - } else { - try { - const consumeResult = await consumeKuaishouEticket({ - baseUrl: eticketSource.baseUrl, - cookie: shopConfig.cookie, - eTicketId: ticketCode, - oid: String(flow.ticket.oid || '').trim(), - formToken: String(flow.ticket.formToken || '').trim(), - }) - - if (consumeResult.consumed) { - consumeStatus = 'success' - consumedAt = now - } else { - consumeStatus = 'failed' - consumeErrorMessage = String(consumeResult.errorMessage || '快手核销失败').trim() - } - } catch (error) { - consumeStatus = 'failed' - consumeErrorMessage = error instanceof Error ? error.message : '快手核销失败' - } - } - - if (consumeStatus !== 'success') { - nextTaskStatus = 'manual_review' - nextResultCode = 'kuaishou_cloud_consume_failed' - nextResultMessage = consumeErrorMessage || '号码已退还,但快手核销未完成,请人工处理' - } - - 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: shopId || flow.consume.shopId, - shopName: String(flow.consume.shopName || order?.shop_name || '').trim(), - 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 === 'success' ? now : task.redeemed_at, - 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) - - await createTaskEvent( - task.id, - consumeStatus === 'success' ? 'kuaishou_cloud_consumed' : 'kuaishou_cloud_consume_failed', - { - ticketCodeMasked: maskCode(ticketCode), - shopId, - consumeStatus, - errorMessage: consumeErrorMessage, - }, - now, - ) - - return { - task: mapTaskActionPayload(updatedTask), - } -} - -/** - * @param {AdminInventoryImportInput} payload - * @returns {Array>>} - */ -function normalizeInventoryImportRows(payload) { - const rows = [] - const directRows = Array.isArray(payload.rows) ? payload.rows : [] - const bulkCodes = Array.isArray(payload.codes) ? payload.codes : [] - - if (directRows.length > 0) { - for (const row of directRows) { - const skuCode = String(row?.skuCode || '').trim() - const displayValue = String(row?.displayValue || '').trim() - const batchNo = String(row?.batchNo || '').trim() - const credentialType = String(row?.credentialType || payload.credentialType || 'tencent_code').trim() || 'tencent_code' - const inventoryGroupCode = String(row?.inventoryGroupCode || payload.inventoryGroupCode || '').trim() - - if (!skuCode || !displayValue) { - continue - } - - rows.push({ skuCode, displayValue, batchNo, credentialType, inventoryGroupCode }) - } - } - - if (bulkCodes.length > 0) { - const skuCode = String(payload.skuCode || '').trim() - const batchNo = String(payload.batchNo || '').trim() - const credentialType = String(payload.credentialType || 'tencent_code').trim() || 'tencent_code' - const inventoryGroupCode = String(payload.inventoryGroupCode || '').trim() - - if (!skuCode) { - throw createHttpError('批量导入时缺少 skuCode', { - statusCode: 400, - errorCode: 'admin_inventory_import_missing_sku', - }) - } - - for (const rawValue of bulkCodes) { - const displayValue = String(rawValue || '').trim() - if (!displayValue) { - continue - } - rows.push({ skuCode, displayValue, batchNo, credentialType, inventoryGroupCode }) - } - } - - return dedupeRows(rows) -} - -function isKuaishouCloudTask(task) { - return String(task?.executor_key || '').trim() === 'kuaishou_ct_assisted' -} - -function resolvePersistedCloudtentaclesContext() { - const source = getCloudtentaclesSourceConfig() - const session = getCloudtentaclesSessionState() - const token = String(session.token || '').trim() - - if (!token) { - throw createHttpError('当前 cloudtentacles 没有可用 token,请先到平台配置完成登录校验', { - statusCode: 409, - errorCode: 'admin_task_kuaishou_cloud_missing_cloud_token', - }) - } - - return { - baseUrl: String(session.baseUrl || source.baseUrl || '').trim() || 'https://123.207.217.176', - token, - deviceId: String(session.deviceId || source.deviceId || '-').trim() || '-', - deviceType: Number(session.deviceType ?? source.deviceType ?? 0), - } -} - -function normalizeKuaishouCloudFlow(value) { - const source = value && typeof value === 'object' ? value : {} - const binding = source.binding && typeof source.binding === 'object' ? source.binding : {} - const role = source.role && typeof source.role === 'object' ? source.role : {} - const purchase = source.purchase && typeof source.purchase === 'object' ? source.purchase : {} - const dispatch = source.dispatch && typeof source.dispatch === 'object' ? source.dispatch : {} - const returnNumber = source.returnNumber && typeof source.returnNumber === 'object' ? source.returnNumber : {} - const consume = source.consume && typeof source.consume === 'object' ? source.consume : {} - const ticket = source.ticket && typeof source.ticket === 'object' ? source.ticket : {} - - return { - ...source, - configId: String(source.configId || '').trim(), - internalSkuCode: String(source.internalSkuCode || '').trim(), - internalSkuName: String(source.internalSkuName || '').trim(), - ticket: { - code: String(ticket.code || '').trim(), - status: String(ticket.status || 'pending').trim() || 'pending', - capturedAt: ticket.capturedAt || null, - capturedBy: ticket.capturedBy || null, - verifiedAt: ticket.verifiedAt || null, - oid: String(ticket.oid || '').trim(), - formToken: String(ticket.formToken || '').trim(), - leftCount: Number(ticket.leftCount || 0) || 0, - goodsTitle: String(ticket.goodsTitle || '').trim(), - }, - binding: { - prepareStatus: String(binding.prepareStatus || 'pending').trim() || 'pending', - cloudSourceKey: String(binding.cloudSourceKey || 'default').trim() || 'default', - skuId: Number(binding.skuId || 0) || 0, - skuName: String(binding.skuName || '').trim(), - vnKey: String(binding.vnKey || '').trim(), - vnId: Number(binding.vnId || 0) || 0, - vnPhone: String(binding.vnPhone || '').trim(), - bindUrl: String(binding.bindUrl || '').trim(), - bindPreparedAt: binding.bindPreparedAt || null, - }, - role: { - status: String(role.status || 'pending').trim() || 'pending', - name: String(role.name || '').trim(), - rid: String(role.rid || '').trim(), - refreshedAt: role.refreshedAt || null, - errorMessage: String(role.errorMessage || '').trim(), - rawInfo: role.rawInfo && typeof role.rawInfo === 'object' ? role.rawInfo : null, - }, - purchase: { - autoBuyEnabled: purchase.autoBuyEnabled !== false, - minAssetReserve: Number(purchase.minAssetReserve || 0) || 0, - usedKnapsack: purchase.usedKnapsack === true, - purchaseTriggered: purchase.purchaseTriggered === true, - assetBefore: Number(purchase.assetBefore || 0) || 0, - assetAfter: Number(purchase.assetAfter || 0) || 0, - purchaseAt: purchase.purchaseAt || null, - }, - dispatch: { - status: String(dispatch.status || 'pending').trim() || 'pending', - dispatchAt: dispatch.dispatchAt || null, - dispatchBy: dispatch.dispatchBy || null, - sendType: Number(dispatch.sendType || 0) || 0, - note: String(dispatch.note || '').trim(), - }, - returnNumber: { - status: String(returnNumber.status || 'pending').trim() || 'pending', - returnedAt: returnNumber.returnedAt || null, - returnedBy: returnNumber.returnedBy || null, - autoReturnEnabled: returnNumber.autoReturnEnabled === true, - }, - consume: { - status: String(consume.status || 'pending').trim() || 'pending', - shopId: String(consume.shopId || '').trim(), - shopName: String(consume.shopName || '').trim(), - autoConsumeEnabled: consume.autoConsumeEnabled === true, - consumedAt: consume.consumedAt || null, - errorMessage: String(consume.errorMessage || '').trim(), - }, - } -} - -function normalizeKuaishouCloudRoleInfo(value) { - const rawInfo = value && typeof value === 'object' ? value : null - - return { - name: String(rawInfo?.name || rawInfo?.roleName || rawInfo?.nickname || '').trim(), - rid: String(rawInfo?.rid || rawInfo?.roleId || rawInfo?.uid || '').trim(), - rawInfo, - } -} - -function resolveKuaishouCloudBindingResources(flow, { skuItems = [], knapsackItems = [] } = {}) { - const normalizedSkuItems = Array.isArray(skuItems) ? skuItems.filter(isCloudSkuLikeItem) : [] - const normalizedKnapsackItems = Array.isArray(knapsackItems) ? knapsackItems.filter(isCloudSkuLikeItem) : [] - const currentSkuId = Number(flow?.binding?.skuId || 0) || 0 - const currentVnKey = String(flow?.binding?.vnKey || '').trim() - const currentSkuName = String(flow?.binding?.skuName || '').trim() - const nameCandidates = collectKuaishouCloudNameCandidates(flow) - - const skuItemById = currentSkuId > 0 - ? normalizedSkuItems.find((item) => Number(item.id || 0) === currentSkuId) || null - : null - const knapsackItemById = currentSkuId > 0 - ? normalizedKnapsackItems.find((item) => Number(item.id || 0) === currentSkuId) || null - : null - - if (skuItemById || knapsackItemById) { - const matchedItem = skuItemById || knapsackItemById - return { - skuId: Number(matchedItem?.id || 0) || 0, - skuName: String(currentSkuName || matchedItem?.name || '').trim(), - vnKey: KUAISHOU_CLOUD_FIXED_VN_KEY, - skuItem: skuItemById, - knapsackItem: knapsackItemById, - resolvedByName: false, - } - } - - const matchedSkuItem = findCloudItemByNames(normalizedSkuItems, nameCandidates) - const matchedKnapsackItem = findCloudItemByNames( - normalizedKnapsackItems, - nameCandidates, - matchedSkuItem ? Number(matchedSkuItem.id || 0) : 0, - ) - const matchedItem = matchedSkuItem || matchedKnapsackItem - - return { - skuId: Number(matchedItem?.id || 0) || 0, - skuName: String(currentSkuName || matchedItem?.name || '').trim(), - vnKey: KUAISHOU_CLOUD_FIXED_VN_KEY, - skuItem: matchedSkuItem, - knapsackItem: matchedKnapsackItem, - resolvedByName: Boolean(matchedItem), - } -} - -/** - * @param {{ flow?: any, binding?: any }} [input] - */ -function resolveKuaishouCloudVnKeyCandidates(input = {}) { - return [KUAISHOU_CLOUD_FIXED_VN_KEY] -} - -/** - * @param {{ cloudContext?: Record, vnKeyCandidates?: string[] }} [input] - */ -async function prepareKuaishouCloudBindResourceWithFallback(input = {}) { - const { cloudContext = {}, vnKeyCandidates = [] } = input - const candidates = Array.isArray(vnKeyCandidates) ? vnKeyCandidates : [] - let lastError = null - - for (const vnKey of candidates) { - let vnId = 0 - let vnPhone = '' - - try { - const appointed = await appointCloudtentaclesVirtualNumber({ - ...cloudContext, - key: vnKey, - }) - vnId = Number(appointed.item?.id || 0) - vnPhone = String(appointed.item?.phone || '').trim() - - if (!vnId || !vnPhone) { - throw createHttpError('申请虚拟号成功但返回数据不完整', { - statusCode: 502, - errorCode: 'admin_task_kuaishou_cloud_invalid_vn', - }) - } - - await generateCloudtentaclesLoginCode({ - ...cloudContext, - key: vnKey, - id: vnId, - }) - - const fetchedCode = await fetchCloudtentaclesVirtualNumberCode({ - ...cloudContext, - key: vnKey, - phone: vnPhone, - }) - - await verifyCloudtentaclesLoginCode({ - ...cloudContext, - key: vnKey, - id: vnId, - code: fetchedCode.code, - }) - - const bindUrlResult = await getCloudtentaclesBindUrl({ - ...cloudContext, - key: vnKey, - id: vnId, - }) - - return { - vnKey, - vnId, - vnPhone, - bindUrl: String(bindUrlResult.bindUrl || '').trim(), - } - } catch (error) { - lastError = error - - if (vnId > 0) { - try { - await backCloudtentaclesVirtualNumber({ - ...cloudContext, - key: vnKey, - id: vnId, - }) - } catch { - // 兜底退号失败时保留原始错误,避免吞掉主链路异常 - } - } - - if (!isRecoverableKuaishouCloudVnKeyError(error)) { - throw error - } - } - } - - throw lastError || createHttpError('没有找到可用的 VN Key', { - statusCode: 409, - errorCode: 'admin_task_kuaishou_cloud_missing_binding_config', - }) -} - -function collectKuaishouCloudNameCandidates(flow) { - return Array.from(new Set([ - String(flow?.binding?.skuName || '').trim(), - String(flow?.internalSkuName || '').trim(), - String(flow?.internalSkuCode || '').trim(), - ].filter(Boolean))) -} - -function findCloudItemByNames(items, nameCandidates, preferredId = 0) { - const normalizedItems = Array.isArray(items) ? items : [] - const normalizedNames = nameCandidates - .map((item) => ({ - raw: String(item || '').trim(), - normalized: normalizeProductName(item), - })) - .filter((item) => item.raw && item.normalized) - - if (normalizedNames.length === 0 || normalizedItems.length === 0) { - return preferredId > 0 - ? normalizedItems.find((item) => Number(item.id || 0) === preferredId) || null - : null - } - - if (preferredId > 0) { - const preferred = normalizedItems.find((item) => Number(item.id || 0) === preferredId) || null - if (preferred) { - return preferred - } - } - - const exactMatches = normalizedItems.filter((item) => { - const itemName = normalizeProductName(item.name) - return normalizedNames.some((candidate) => candidate.normalized === itemName) - }) - if (exactMatches.length > 0) { - return exactMatches[0] - } - - const partialMatches = normalizedItems.filter((item) => { - const itemName = normalizeProductName(item.name) - return normalizedNames.some((candidate) => itemName.includes(candidate.normalized) || candidate.normalized.includes(itemName)) - }) - if (partialMatches.length > 0) { - return partialMatches.sort((left, right) => String(left.name || '').length - String(right.name || '').length)[0] - } - - return null -} - -function isCloudSkuLikeItem(item) { - return Boolean(item) && typeof item === 'object' && Number(item.id || 0) > 0 -} - -function maskPhone(value) { - const text = String(value || '').trim() - if (!text) { - return '' - } - - if (text.length <= 7) { - return `${text.slice(0, 2)}***${text.slice(-2)}` - } - - return `${text.slice(0, 3)}****${text.slice(-4)}` -} - -function maskCode(value) { - const text = String(value || '').trim() - if (!text) { - return '' - } - - if (text.length <= 8) { - return `${text.slice(0, 2)}***${text.slice(-2)}` - } - - return `${text.slice(0, 4)}****${text.slice(-4)}` -} - -/** - * @param {Array>>} rows - */ -function dedupeRows(rows) { - const seen = new Set() - const output = [] - - for (const row of rows) { - const key = `${row.skuCode}::${row.credentialType || 'tencent_code'}::${row.displayValue}` - if (seen.has(key)) { - continue - } - seen.add(key) - output.push(row) - } - - return output -} - -function resolveTaskInventoryGroupCodes(task) { - const taskContext = parseTaskContext(task) - const inventoryGroupCode = String(taskContext?.inventoryGroupCode || '').trim() - return inventoryGroupCode ? [inventoryGroupCode] : null -} - - -/** - * @param {Awaited>} task - * @param {ReturnType} viewerContext - * @param {'confirm' | 'redeem'} action - */ -function ensureViewerCanOperateAssistedTask(task, viewerContext, action) { - if (!viewerContext.canOperateAssistedTask || !isAssistedClaimTask(task)) { - throw createHttpError('当前账号没有此操作权限', { - statusCode: 403, - errorCode: 'admin_task_assisted_permission_denied', - }) - } - - if (action === 'confirm' && !canViewerConfirmAssistedRole(task, viewerContext)) { - throw createHttpError('当前任务状态还不能确认角色', { - statusCode: 409, - errorCode: 'admin_task_assisted_confirm_not_allowed', - }) - } - - if (action === 'redeem' && !canViewerRedeemAssistedTask(task, viewerContext)) { - throw createHttpError('当前任务状态还不能开始兑换', { - statusCode: 409, - errorCode: 'admin_task_assisted_redeem_not_allowed', - }) - } -} - -function isRecoverableTaskSessionCloseError(error) { - const errorCode = String(error?.errorCode || error?.code || '').trim() - return errorCode === 'session_not_found' || errorCode === 'session_closed' -} - -function isRecoverableKuaishouCloudVnKeyError(error) { - const errorCode = String(error?.errorCode || error?.code || '').trim() - const errorMessage = String(error?.message || '').trim() - return errorCode === 'cloudtentacles_vn_bind_url_failed' - && errorMessage.includes('不支持的游戏类型') -} - -function normalizeManualDispatchOutcome(value) { - const normalized = String(value || '').trim().toLowerCase() - - if (normalized === 'failed') { - return 'failed' - } - - return 'delivered' -} - -function getTaskClaimExpiresAt(task) { - return task?.claim_expires_at || task?.primary_claim_expires_at || null -} - -async function ensureTaskClaimLink(task) { - const tokenStatus = String(task?.primary_claim_token_status || '').trim() - const token = String(task?.primary_claim_token || task?.claim_token || '').trim() - const expiredAt = getTaskClaimExpiresAt(task) - - if (tokenStatus === 'active' && token && !isClaimExpired(expiredAt)) { - return { - token, - expiredAt, - claimUrl: buildClaimUrl(token), - } - } - - const claimToken = await createTaskClaimToken(task.id) - return { - token: claimToken.token, - expiredAt: claimToken.expired_at, - claimUrl: claimToken.claimUrl, - } -} - -function isClaimExpired(expiredAt) { - if (!expiredAt) { - return false - } - - const timestamp = new Date(expiredAt).getTime() - return Number.isFinite(timestamp) && timestamp <= Date.now() -} +export { + createAdminInventoryItem, + importAdminInventoryItems, + releaseAdminInventoryItem, + invalidateAdminInventoryItem, +} from './write/inventory.js' + +export { + replayAdminWebhookEvent, +} from './write/webhook-events.js' + +export { + releaseAdminTaskInventory, + releaseAdminTaskInventoryBinding, + regenerateAdminTaskClaimLink, + confirmAdminTaskAssistedRole, + redeemAdminTaskAssisted, + closeAdminTask, + closeAdminTaskWithDeps, + markAdminTaskManualReview, + retryAdminTask, + completeAdminTaskManualDispatch, +} from './write/task-actions.js' + +export { + prepareAdminTaskKuaishouCloudFulfillment, + dispatchAdminTaskKuaishouCloudFulfillment, + refreshAdminTaskKuaishouCloudRoleInfo, + returnNumberAdminTaskKuaishouCloudFulfillment, +} from './write/kuaishou-cloud-actions.js' diff --git a/apps/backend/src/services/admin/write/inventory.js b/apps/backend/src/services/admin/write/inventory.js new file mode 100644 index 00000000..bd0ff0ab --- /dev/null +++ b/apps/backend/src/services/admin/write/inventory.js @@ -0,0 +1,215 @@ +// @ts-check + +import { + createInventoryItems, + invalidateInventoryItem, + listInventoryItems, + releaseReservedInventoryItem, +} from '../../../repositories/inventory-repo.js' +import { createHttpError } from '../../../utils/http.js' +import { nowIso } from '../../../utils/time.js' +import { + getRequiredInventoryItem, + mapAdminInventoryListItem, +} from '../admin-inventory-read-helpers.js' + +/** @typedef {import('../../../types/admin-read-inputs.js').AdminEntityIdInput} AdminEntityIdInput */ +/** @typedef {import('../../../types/admin-write-inputs.js').AdminInventoryCreateInput} AdminInventoryCreateInput */ +/** @typedef {import('../../../types/admin-write-inputs.js').AdminInventoryImportInput} AdminInventoryImportInput */ +/** @typedef {import('../../../types/admin-write-inputs.js').AdminInventoryImportRowInput} AdminInventoryImportRowInput */ +/** @typedef {import('../../../types/admin-write-inputs.js').AdminInventoryInvalidateInput} AdminInventoryInvalidateInput */ +/** @typedef {import('../../../types/admin-write-models.js').AdminInventoryImportResponse} AdminInventoryImportResponse */ +/** @typedef {import('../../../types/admin-write-models.js').AdminInventoryMutationResponse} AdminInventoryMutationResponse */ + +/** @returns {Promise} */ +/** @param {AdminInventoryCreateInput} [payload] */ +export async function createAdminInventoryItem(payload = /** @type {AdminInventoryCreateInput} */ ({})) { + const skuCode = String(payload.skuCode || '').trim() + const displayValue = String(payload.displayValue || '').trim() + const batchNo = String(payload.batchNo || '').trim() + const credentialType = String(payload.credentialType || 'tencent_code').trim() || 'tencent_code' + const inventoryGroupCode = String(payload.inventoryGroupCode || '').trim() + + if (!skuCode || !displayValue) { + throw createHttpError('缺少 skuCode 或 displayValue', { + statusCode: 400, + errorCode: 'admin_inventory_create_invalid', + }) + } + + const now = nowIso() + const created = await createInventoryItems([ + { + skuCode, + displayValue, + batchNo, + credentialType, + inventoryGroupCode, + createdAt: now, + updatedAt: now, + }, + ]) + + if (created === 0) { + throw createHttpError('库存凭据已存在,不能重复新增', { + statusCode: 409, + errorCode: 'admin_inventory_duplicate', + }) + } + + const { items } = await listInventoryItems({ + page: 1, + pageSize: 20, + skuCode, + inventoryGroupCode, + }) + const createdItem = items.find((item) => item.display_value === displayValue && item.inventory_group_code === inventoryGroupCode) || null + + return { + inventoryItem: createdItem ? await mapAdminInventoryListItem(createdItem) : null, + } +} + +/** @returns {Promise} */ +/** @param {AdminInventoryImportInput} [payload] */ +export async function importAdminInventoryItems(payload = /** @type {AdminInventoryImportInput} */ ({})) { + const rows = normalizeInventoryImportRows(payload) + + if (rows.length === 0) { + throw createHttpError('没有可导入的库存凭据数据', { + statusCode: 400, + errorCode: 'admin_inventory_import_empty', + }) + } + + const now = nowIso() + const normalizedRows = rows.map((row) => ({ + skuCode: row.skuCode, + batchNo: row.batchNo, + displayValue: row.displayValue, + credentialType: row.credentialType, + inventoryGroupCode: row.inventoryGroupCode, + createdAt: now, + updatedAt: now, + })) + + const created = await createInventoryItems(normalizedRows) + + return { + total: normalizedRows.length, + created, + duplicated: normalizedRows.length - created, + } +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} inventoryItemId */ +export async function releaseAdminInventoryItem(inventoryItemId) { + const inventoryItem = await getRequiredInventoryItem(inventoryItemId) + + if (inventoryItem.status !== 'reserved') { + throw createHttpError('当前库存项不是预占状态,不能释放', { + statusCode: 409, + errorCode: 'admin_inventory_release_not_allowed', + }) + } + + const updated = await releaseReservedInventoryItem(inventoryItem.id, nowIso()) + + return { + inventoryItem: await mapAdminInventoryListItem(updated), + } +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} inventoryItemId */ +/** @param {AdminInventoryInvalidateInput} [payload] */ +export async function invalidateAdminInventoryItem( + inventoryItemId, + payload = /** @type {AdminInventoryInvalidateInput} */ ({}), +) { + const inventoryItem = await getRequiredInventoryItem(inventoryItemId) + + if (inventoryItem.status !== 'available') { + throw createHttpError('只有可用库存项才能作废,请先释放预占', { + statusCode: 409, + errorCode: 'admin_inventory_invalidate_not_allowed', + }) + } + + const reason = String(payload.reason || '').trim() || '后台手动作废' + const updated = await invalidateInventoryItem(inventoryItem.id, reason, nowIso()) + + return { + inventoryItem: await mapAdminInventoryListItem(updated), + } +} + +/** + * @param {AdminInventoryImportInput} payload + * @returns {Array>>} + */ +function normalizeInventoryImportRows(payload) { + const rows = [] + const directRows = Array.isArray(payload.rows) ? payload.rows : [] + const bulkCodes = Array.isArray(payload.codes) ? payload.codes : [] + + if (directRows.length > 0) { + for (const row of directRows) { + const skuCode = String(row?.skuCode || '').trim() + const displayValue = String(row?.displayValue || '').trim() + const batchNo = String(row?.batchNo || '').trim() + const credentialType = String(row?.credentialType || payload.credentialType || 'tencent_code').trim() || 'tencent_code' + const inventoryGroupCode = String(row?.inventoryGroupCode || payload.inventoryGroupCode || '').trim() + + if (!skuCode || !displayValue) { + continue + } + + rows.push({ skuCode, displayValue, batchNo, credentialType, inventoryGroupCode }) + } + } + + if (bulkCodes.length > 0) { + const skuCode = String(payload.skuCode || '').trim() + const batchNo = String(payload.batchNo || '').trim() + const credentialType = String(payload.credentialType || 'tencent_code').trim() || 'tencent_code' + const inventoryGroupCode = String(payload.inventoryGroupCode || '').trim() + + if (!skuCode) { + throw createHttpError('批量导入时缺少 skuCode', { + statusCode: 400, + errorCode: 'admin_inventory_import_missing_sku', + }) + } + + for (const rawValue of bulkCodes) { + const displayValue = String(rawValue || '').trim() + if (!displayValue) { + continue + } + rows.push({ skuCode, displayValue, batchNo, credentialType, inventoryGroupCode }) + } + } + + return dedupeRows(rows) +} + +/** + * @param {Array>>} rows + */ +function dedupeRows(rows) { + const seen = new Set() + const output = [] + + for (const row of rows) { + const key = `${row.skuCode}::${row.credentialType || 'tencent_code'}::${row.displayValue}` + if (seen.has(key)) { + continue + } + seen.add(key) + output.push(row) + } + + return output +} diff --git a/apps/backend/src/services/admin/write/kuaishou-cloud-actions.js b/apps/backend/src/services/admin/write/kuaishou-cloud-actions.js new file mode 100644 index 00000000..74c44e50 --- /dev/null +++ b/apps/backend/src/services/admin/write/kuaishou-cloud-actions.js @@ -0,0 +1,567 @@ +// @ts-check + +import { getOrderById } from '../../../repositories/order-repo.js' +import { updateTask } from '../../../repositories/task-repo.js' +import { createTaskEvent } from '../../../repositories/task-event-repo.js' +import { + buyCloudtentaclesSku, + getCloudtentaclesAsset, + listCloudtentaclesSku, + useCloudtentaclesSku, +} from '../../platforms/cloudtentacles/catalog-service.js' +import { getCloudtentaclesKnapsack } from '../../platforms/cloudtentacles/knapsack-service.js' +import { + backCloudtentaclesVirtualNumber, + getCloudtentaclesBindInfo, +} from '../../platforms/cloudtentacles/virtual-number-service.js' +import { consumeKuaishouEticket } from '../../platforms/kuaishou-eticket/consume-service.js' +import { + getKuaishouEticketSourceConfig, + resolveKuaishouEticketShopConfig, +} from '../../platforms/kuaishou-eticket/source-config-service.js' +import { createHttpError } from '../../../utils/http.js' +import { nowIso } from '../../../utils/time.js' +import { + createAdminViewerContext, + parseTaskContext, +} from '../admin-read-shared-helpers.js' +import { + getRequiredTask, + mapTaskActionPayload, +} from '../admin-task-read-helpers.js' +import { + ensureTaskClaimLink, + getTaskClaimExpiresAt, + maskCode, + maskPhone, +} from './shared.js' +import { + isKuaishouCloudTask, + normalizeKuaishouCloudFlow, + normalizeKuaishouCloudRoleInfo, + prepareKuaishouCloudBindResourceWithFallback, + resolveKuaishouCloudBindingResources, + resolveKuaishouCloudVnKeyCandidates, + resolvePersistedCloudtentaclesContext, +} from './kuaishou-cloud-helpers.js' + +/** @typedef {import('../../../types/admin-read-inputs.js').AdminEntityIdInput} AdminEntityIdInput */ +/** @typedef {import('../../../types/admin-read-inputs.js').AdminViewerSessionInput} AdminViewerSessionInput */ +/** @typedef {import('../../../types/admin-write-inputs.js').AdminTaskKuaishouCloudDispatchInput} AdminTaskKuaishouCloudDispatchInput */ +/** @typedef {import('../../../types/admin-write-models.js').AdminTaskActionResponse} AdminTaskActionResponse */ + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +/** @param {AdminViewerSessionInput | null} [session] */ +export async function prepareAdminTaskKuaishouCloudFulfillment(taskId, session = null) { + const task = await getRequiredTask(taskId) + const now = nowIso() + const viewerContext = createAdminViewerContext(session) + + if (!viewerContext.canManageTaskLifecycle) { + throw createHttpError('当前账号没有权限准备绑定资源', { + statusCode: 403, + errorCode: 'admin_task_prepare_kuaishou_cloud_forbidden', + }) + } + + if (!isKuaishouCloudTask(task)) { + throw createHttpError('当前任务不是快手 cloud 履约任务', { + statusCode: 409, + errorCode: 'admin_task_not_kuaishou_cloud', + }) + } + + const taskContext = parseTaskContext(task) + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) + const cloudContext = resolvePersistedCloudtentaclesContext() + const [knapsack, skuList] = await Promise.all([ + getCloudtentaclesKnapsack(cloudContext), + listCloudtentaclesSku(cloudContext), + ]) + const resolvedBinding = resolveKuaishouCloudBindingResources(flow, { + skuItems: Array.isArray(skuList.items) ? skuList.items : [], + knapsackItems: Array.isArray(knapsack.items) ? knapsack.items : [], + }) + const vnKeyCandidates = resolveKuaishouCloudVnKeyCandidates({ + flow, + binding: resolvedBinding, + }) + + if (!resolvedBinding.skuId || vnKeyCandidates.length === 0) { + throw createHttpError('当前任务缺少可用 cloud 资源,且无法根据商品名自动解析 SKU / VN Key', { + statusCode: 409, + errorCode: 'admin_task_kuaishou_cloud_missing_binding_config', + }) + } + + const flowWithResolvedBinding = { + ...flow, + binding: { + ...flow.binding, + skuId: resolvedBinding.skuId, + skuName: resolvedBinding.skuName, + vnKey: resolvedBinding.vnKey, + }, + } + const knapsackItem = resolvedBinding.knapsackItem + let usedKnapsack = Number(knapsackItem?.count || 0) > 0 + let purchaseTriggered = false + let assetBefore = 0 + let assetAfter = 0 + const claimLinkState = await ensureTaskClaimLink(task) + + if (!usedKnapsack) { + if (!flowWithResolvedBinding.purchase.autoBuyEnabled) { + throw createHttpError('背包中没有现成库存,且当前配置未开启自动购买', { + statusCode: 409, + errorCode: 'admin_task_kuaishou_cloud_auto_buy_disabled', + }) + } + + const asset = await getCloudtentaclesAsset(cloudContext) + const targetSku = resolvedBinding.skuItem + + if (!targetSku) { + throw createHttpError(`cloudtentacles 未找到 SKU ${flowWithResolvedBinding.binding.skuId}`, { + statusCode: 404, + errorCode: 'admin_task_kuaishou_cloud_sku_not_found', + }) + } + + assetBefore = Number(asset.asset || 0) || 0 + const targetPrice = Number(targetSku.price || 0) || 0 + const requiredAsset = targetPrice + flowWithResolvedBinding.purchase.minAssetReserve + if (assetBefore < requiredAsset) { + throw createHttpError(`余额不足,当前 ${assetBefore},至少需要 ${requiredAsset}`, { + statusCode: 409, + errorCode: 'admin_task_kuaishou_cloud_asset_not_enough', + }) + } + + await buyCloudtentaclesSku({ + ...cloudContext, + id: flowWithResolvedBinding.binding.skuId, + count: 1, + }) + purchaseTriggered = true + + const assetResult = await getCloudtentaclesAsset(cloudContext) + assetAfter = Number(assetResult.asset || 0) || 0 + } + + const preparedBinding = await prepareKuaishouCloudBindResourceWithFallback({ + cloudContext, + vnKeyCandidates, + }) + const vnId = preparedBinding.vnId + const vnPhone = preparedBinding.vnPhone + + const nextContext = { + ...taskContext, + kuaishouCloudFulfillment: { + ...flowWithResolvedBinding, + binding: { + ...flowWithResolvedBinding.binding, + vnKey: preparedBinding.vnKey, + prepareStatus: 'ready', + vnId, + vnPhone, + bindUrl: preparedBinding.bindUrl, + bindPreparedAt: now, + roleName: '', + roleId: '', + }, + role: { + status: 'pending', + name: '', + rid: '', + refreshedAt: null, + errorMessage: '', + rawInfo: null, + }, + purchase: { + ...flowWithResolvedBinding.purchase, + usedKnapsack, + purchaseTriggered, + assetBefore, + assetAfter, + purchaseAt: purchaseTriggered ? now : flowWithResolvedBinding.purchase.purchaseAt, + }, + }, + } + + const updatedTask = await updateTask(task.id, { + task_status: 'waiting_binding', + inventory_status: 'not_required', + user_action_status: 'pending', + claim_token: claimLinkState.token || task.claim_token || '', + claim_expires_at: claimLinkState.expiredAt || getTaskClaimExpiresAt(task), + last_error: '', + context_json: JSON.stringify(nextContext), + updated_at: now, + }) + + await createTaskEvent(task.id, 'kuaishou_cloud_binding_prepared', { + skuId: flowWithResolvedBinding.binding.skuId, + skuName: flowWithResolvedBinding.binding.skuName, + vnKey: preparedBinding.vnKey, + vnId, + vnPhoneMasked: maskPhone(vnPhone), + bindUrl: preparedBinding.bindUrl, + purchaseTriggered, + usedKnapsack, + resolvedByName: resolvedBinding.resolvedByName, + }, now) + + return { + task: mapTaskActionPayload(updatedTask), + claimUrl: claimLinkState.claimUrl, + token: claimLinkState.token, + } +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +/** @param {AdminTaskKuaishouCloudDispatchInput} [payload] */ +/** @param {AdminViewerSessionInput | null} [session] */ +export async function dispatchAdminTaskKuaishouCloudFulfillment(taskId, payload = {}, session = null) { + const task = await getRequiredTask(taskId) + const now = nowIso() + const viewerContext = createAdminViewerContext(session) + + if (!viewerContext.canManageTaskLifecycle) { + throw createHttpError('当前账号没有权限执行发货', { + statusCode: 403, + errorCode: 'admin_task_dispatch_kuaishou_cloud_forbidden', + }) + } + + if (!isKuaishouCloudTask(task)) { + throw createHttpError('当前任务不是快手 cloud 履约任务', { + statusCode: 409, + errorCode: 'admin_task_not_kuaishou_cloud', + }) + } + + const taskContext = parseTaskContext(task) + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) + const cloudContext = resolvePersistedCloudtentaclesContext() + + if (!flow.binding.skuId || !flow.binding.vnId || !flow.binding.vnPhone) { + throw createHttpError('当前任务还没有准备好绑定资源,请先准备绑定资源', { + statusCode: 409, + errorCode: 'admin_task_kuaishou_cloud_not_prepared', + }) + } + + const ticketCode = String(payload.ticketCode || '').trim() + const persistedTicketCode = String(flow.ticket.code || '').trim() + + if (!persistedTicketCode && !ticketCode) { + throw createHttpError('客户还没有在领取页提交核销码,暂时不能直接发货', { + statusCode: 409, + errorCode: 'admin_task_kuaishou_cloud_missing_ticket_code', + }) + } + + const dispatchResult = await useCloudtentaclesSku({ + ...cloudContext, + id: flow.binding.skuId, + virtualNumberId: flow.binding.vnId, + phone: flow.binding.vnPhone, + }) + + const nextContext = { + ...taskContext, + kuaishouCloudFulfillment: { + ...flow, + ticket: { + ...flow.ticket, + code: ticketCode || persistedTicketCode, + capturedAt: ticketCode ? now : flow.ticket.capturedAt, + capturedBy: ticketCode && session + ? { + userId: Number(session.userId || 0) || 0, + username: String(session.username || '').trim(), + role: String(session.role || '').trim(), + } + : flow.ticket.capturedBy, + }, + dispatch: { + ...flow.dispatch, + status: 'success', + dispatchAt: now, + dispatchBy: session + ? { + userId: Number(session.userId || 0) || 0, + username: String(session.username || '').trim(), + role: String(session.role || '').trim(), + } + : null, + sendType: Number(dispatchResult.sendType || 0) || 0, + note: String(dispatchResult.note || dispatchResult.responseMessage || '').trim(), + }, + }, + } + + const updatedTask = await updateTask(task.id, { + task_status: 'dispatched_pending_return', + delivery_status: 'delivered', + result_code: 'kuaishou_cloud_dispatched', + result_message: String(dispatchResult.responseMessage || dispatchResult.note || 'cloudtentacles 发货成功').trim(), + last_error: '', + context_json: JSON.stringify(nextContext), + updated_at: now, + }) + + await createTaskEvent(task.id, 'kuaishou_cloud_dispatched', { + ticketCodeMasked: maskCode(ticketCode || persistedTicketCode), + skuId: flow.binding.skuId, + vnId: flow.binding.vnId, + vnPhoneMasked: maskPhone(flow.binding.vnPhone), + sendType: dispatchResult.sendType, + note: dispatchResult.note, + }, now) + + return { + task: mapTaskActionPayload(updatedTask), + } +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +/** @param {AdminViewerSessionInput | null} [session] */ +export async function refreshAdminTaskKuaishouCloudRoleInfo(taskId, session = null) { + const task = await getRequiredTask(taskId) + const now = nowIso() + const viewerContext = createAdminViewerContext(session) + + if (!viewerContext.canOperateAssistedTask) { + throw createHttpError('当前账号没有权限刷新角色信息', { + statusCode: 403, + errorCode: 'admin_task_refresh_kuaishou_cloud_role_forbidden', + }) + } + + if (!isKuaishouCloudTask(task)) { + throw createHttpError('当前任务不是快手 cloud 履约任务', { + statusCode: 409, + errorCode: 'admin_task_not_kuaishou_cloud', + }) + } + + const taskContext = parseTaskContext(task) + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) + if (!flow.binding.vnId || !flow.binding.vnKey) { + throw createHttpError('当前任务还没有可查询的绑定角色信息,请先准备绑定资源', { + statusCode: 409, + errorCode: 'admin_task_kuaishou_cloud_missing_bind_info_context', + }) + } + + const cloudContext = resolvePersistedCloudtentaclesContext() + 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), + } +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +/** @param {AdminViewerSessionInput | null} [session] */ +export async function returnNumberAdminTaskKuaishouCloudFulfillment(taskId, session = null) { + const task = await getRequiredTask(taskId) + const now = nowIso() + const viewerContext = createAdminViewerContext(session) + + if (!viewerContext.canManageTaskLifecycle) { + throw createHttpError('当前账号没有权限退还号码', { + statusCode: 403, + errorCode: 'admin_task_return_kuaishou_cloud_forbidden', + }) + } + + if (!isKuaishouCloudTask(task)) { + throw createHttpError('当前任务不是快手 cloud 履约任务', { + statusCode: 409, + errorCode: 'admin_task_not_kuaishou_cloud', + }) + } + + const taskContext = parseTaskContext(task) + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) + const cloudContext = resolvePersistedCloudtentaclesContext() + + if (!flow.binding.vnId || !flow.binding.vnKey) { + throw createHttpError('当前任务缺少可退还的虚拟号信息', { + statusCode: 409, + errorCode: 'admin_task_kuaishou_cloud_missing_return_context', + }) + } + + await backCloudtentaclesVirtualNumber({ + ...cloudContext, + key: flow.binding.vnKey, + id: flow.binding.vnId, + }) + + const order = await getOrderById(task.order_id) + const ticketCode = String(flow.ticket.code || '').trim() + const shopId = String(order?.shop_id || '').trim() + const eticketSource = getKuaishouEticketSourceConfig() + const shopConfig = resolveKuaishouEticketShopConfig({ + shopId, + shopName: String(order?.shop_name || '').trim(), + }) + let consumeStatus = 'pending' + let consumeErrorMessage = '' + let consumedAt = null + let nextTaskStatus = 'completed' + let nextResultCode = 'kuaishou_cloud_completed' + let nextResultMessage = 'cloudtentacles 发货、退号并完成快手核销' + + if (!order) { + consumeStatus = 'failed' + consumeErrorMessage = '任务关联订单不存在,无法执行快手核销' + } else if (!ticketCode) { + consumeStatus = 'failed' + consumeErrorMessage = '客户未提交有效核销码,无法执行快手核销' + } else if (!shopConfig || shopConfig.enabled === false || !String(shopConfig.cookie || '').trim()) { + consumeStatus = 'failed' + consumeErrorMessage = '订单对应快手小店缺少可用 Cookie,无法执行快手核销' + } else { + try { + const consumeResult = await consumeKuaishouEticket({ + baseUrl: eticketSource.baseUrl, + cookie: shopConfig.cookie, + eTicketId: ticketCode, + oid: String(flow.ticket.oid || '').trim(), + formToken: String(flow.ticket.formToken || '').trim(), + }) + + if (consumeResult.consumed) { + consumeStatus = 'success' + consumedAt = now + } else { + consumeStatus = 'failed' + consumeErrorMessage = String(consumeResult.errorMessage || '快手核销失败').trim() + } + } catch (error) { + consumeStatus = 'failed' + consumeErrorMessage = error instanceof Error ? error.message : '快手核销失败' + } + } + + if (consumeStatus !== 'success') { + nextTaskStatus = 'manual_review' + nextResultCode = 'kuaishou_cloud_consume_failed' + nextResultMessage = consumeErrorMessage || '号码已退还,但快手核销未完成,请人工处理' + } + + 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: shopId || flow.consume.shopId, + shopName: String(flow.consume.shopName || order?.shop_name || '').trim(), + 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 === 'success' ? now : task.redeemed_at, + 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) + + await createTaskEvent( + task.id, + consumeStatus === 'success' ? 'kuaishou_cloud_consumed' : 'kuaishou_cloud_consume_failed', + { + ticketCodeMasked: maskCode(ticketCode), + shopId, + consumeStatus, + errorMessage: consumeErrorMessage, + }, + now, + ) + + return { + task: mapTaskActionPayload(updatedTask), + } +} diff --git a/apps/backend/src/services/admin/write/kuaishou-cloud-helpers.js b/apps/backend/src/services/admin/write/kuaishou-cloud-helpers.js new file mode 100644 index 00000000..d3b3e50c --- /dev/null +++ b/apps/backend/src/services/admin/write/kuaishou-cloud-helpers.js @@ -0,0 +1,326 @@ +// @ts-check + +import { normalizeProductName } from '../../order/product-match-service.js' +import { getCloudtentaclesSourceConfig } from '../../platforms/cloudtentacles/source-config-service.js' +import { getCloudtentaclesSessionState } from '../../platforms/cloudtentacles/session-state-service.js' +import { + appointCloudtentaclesVirtualNumber, + backCloudtentaclesVirtualNumber, + fetchCloudtentaclesVirtualNumberCode, + generateCloudtentaclesLoginCode, + getCloudtentaclesBindUrl, + verifyCloudtentaclesLoginCode, +} from '../../platforms/cloudtentacles/virtual-number-service.js' +import { createHttpError } from '../../../utils/http.js' + +export const KUAISHOU_CLOUD_FIXED_VN_KEY = '1' + +export function isKuaishouCloudTask(task) { + return String(task?.executor_key || '').trim() === 'kuaishou_ct_assisted' +} + +export function resolvePersistedCloudtentaclesContext() { + const source = getCloudtentaclesSourceConfig() + const session = getCloudtentaclesSessionState() + const token = String(session.token || '').trim() + + if (!token) { + throw createHttpError('当前 cloudtentacles 没有可用 token,请先到平台配置完成登录校验', { + statusCode: 409, + errorCode: 'admin_task_kuaishou_cloud_missing_cloud_token', + }) + } + + return { + baseUrl: String(session.baseUrl || source.baseUrl || '').trim() || 'https://123.207.217.176', + token, + deviceId: String(session.deviceId || source.deviceId || '-').trim() || '-', + deviceType: Number(session.deviceType ?? source.deviceType ?? 0), + } +} + +export function normalizeKuaishouCloudFlow(value) { + const source = value && typeof value === 'object' ? value : {} + const binding = source.binding && typeof source.binding === 'object' ? source.binding : {} + const role = source.role && typeof source.role === 'object' ? source.role : {} + const purchase = source.purchase && typeof source.purchase === 'object' ? source.purchase : {} + const dispatch = source.dispatch && typeof source.dispatch === 'object' ? source.dispatch : {} + const returnNumber = source.returnNumber && typeof source.returnNumber === 'object' ? source.returnNumber : {} + const consume = source.consume && typeof source.consume === 'object' ? source.consume : {} + const ticket = source.ticket && typeof source.ticket === 'object' ? source.ticket : {} + + return { + ...source, + configId: String(source.configId || '').trim(), + internalSkuCode: String(source.internalSkuCode || '').trim(), + internalSkuName: String(source.internalSkuName || '').trim(), + ticket: { + code: String(ticket.code || '').trim(), + status: String(ticket.status || 'pending').trim() || 'pending', + capturedAt: ticket.capturedAt || null, + capturedBy: ticket.capturedBy || null, + verifiedAt: ticket.verifiedAt || null, + oid: String(ticket.oid || '').trim(), + formToken: String(ticket.formToken || '').trim(), + leftCount: Number(ticket.leftCount || 0) || 0, + goodsTitle: String(ticket.goodsTitle || '').trim(), + }, + binding: { + prepareStatus: String(binding.prepareStatus || 'pending').trim() || 'pending', + cloudSourceKey: String(binding.cloudSourceKey || 'default').trim() || 'default', + skuId: Number(binding.skuId || 0) || 0, + skuName: String(binding.skuName || '').trim(), + vnKey: String(binding.vnKey || '').trim(), + vnId: Number(binding.vnId || 0) || 0, + vnPhone: String(binding.vnPhone || '').trim(), + bindUrl: String(binding.bindUrl || '').trim(), + bindPreparedAt: binding.bindPreparedAt || null, + roleName: String(binding.roleName || '').trim(), + roleId: String(binding.roleId || '').trim(), + }, + role: { + status: String(role.status || 'pending').trim() || 'pending', + name: String(role.name || '').trim(), + rid: String(role.rid || '').trim(), + refreshedAt: role.refreshedAt || null, + errorMessage: String(role.errorMessage || '').trim(), + rawInfo: role.rawInfo && typeof role.rawInfo === 'object' ? role.rawInfo : null, + }, + purchase: { + autoBuyEnabled: purchase.autoBuyEnabled !== false, + minAssetReserve: Number(purchase.minAssetReserve || 0) || 0, + usedKnapsack: purchase.usedKnapsack === true, + purchaseTriggered: purchase.purchaseTriggered === true, + assetBefore: Number(purchase.assetBefore || 0) || 0, + assetAfter: Number(purchase.assetAfter || 0) || 0, + purchaseAt: purchase.purchaseAt || null, + }, + dispatch: { + status: String(dispatch.status || 'pending').trim() || 'pending', + dispatchAt: dispatch.dispatchAt || null, + dispatchBy: dispatch.dispatchBy || null, + sendType: Number(dispatch.sendType || 0) || 0, + note: String(dispatch.note || '').trim(), + }, + returnNumber: { + status: String(returnNumber.status || 'pending').trim() || 'pending', + returnedAt: returnNumber.returnedAt || null, + returnedBy: returnNumber.returnedBy || null, + autoReturnEnabled: returnNumber.autoReturnEnabled === true, + }, + consume: { + status: String(consume.status || 'pending').trim() || 'pending', + shopId: String(consume.shopId || '').trim(), + shopName: String(consume.shopName || '').trim(), + autoConsumeEnabled: consume.autoConsumeEnabled === true, + consumedAt: consume.consumedAt || null, + errorMessage: String(consume.errorMessage || '').trim(), + }, + } +} + +export function normalizeKuaishouCloudRoleInfo(value) { + const rawInfo = value && typeof value === 'object' ? value : null + + return { + name: String(rawInfo?.name || rawInfo?.roleName || rawInfo?.nickname || '').trim(), + rid: String(rawInfo?.rid || rawInfo?.roleId || rawInfo?.uid || '').trim(), + rawInfo, + } +} + +export function resolveKuaishouCloudBindingResources(flow, { skuItems = [], knapsackItems = [] } = {}) { + const normalizedSkuItems = Array.isArray(skuItems) ? skuItems.filter(isCloudSkuLikeItem) : [] + const normalizedKnapsackItems = Array.isArray(knapsackItems) ? knapsackItems.filter(isCloudSkuLikeItem) : [] + const currentSkuId = Number(flow?.binding?.skuId || 0) || 0 + const currentSkuName = String(flow?.binding?.skuName || '').trim() + const nameCandidates = collectKuaishouCloudNameCandidates(flow) + + const skuItemById = currentSkuId > 0 + ? normalizedSkuItems.find((item) => Number(item.id || 0) === currentSkuId) || null + : null + const knapsackItemById = currentSkuId > 0 + ? normalizedKnapsackItems.find((item) => Number(item.id || 0) === currentSkuId) || null + : null + + if (skuItemById || knapsackItemById) { + const matchedItem = skuItemById || knapsackItemById + return { + skuId: Number(matchedItem?.id || 0) || 0, + skuName: String(currentSkuName || matchedItem?.name || '').trim(), + vnKey: KUAISHOU_CLOUD_FIXED_VN_KEY, + skuItem: skuItemById, + knapsackItem: knapsackItemById, + resolvedByName: false, + } + } + + const matchedSkuItem = findCloudItemByNames(normalizedSkuItems, nameCandidates) + const matchedKnapsackItem = findCloudItemByNames( + normalizedKnapsackItems, + nameCandidates, + matchedSkuItem ? Number(matchedSkuItem.id || 0) : 0, + ) + const matchedItem = matchedSkuItem || matchedKnapsackItem + + return { + skuId: Number(matchedItem?.id || 0) || 0, + skuName: String(currentSkuName || matchedItem?.name || '').trim(), + vnKey: KUAISHOU_CLOUD_FIXED_VN_KEY, + skuItem: matchedSkuItem, + knapsackItem: matchedKnapsackItem, + resolvedByName: Boolean(matchedItem), + } +} + +/** + * @param {{ flow?: any, binding?: any }} [input] + */ +export function resolveKuaishouCloudVnKeyCandidates(input = {}) { + return [KUAISHOU_CLOUD_FIXED_VN_KEY] +} + +/** + * @param {{ cloudContext?: Record, vnKeyCandidates?: string[] }} [input] + */ +export async function prepareKuaishouCloudBindResourceWithFallback(input = {}) { + const { cloudContext = {}, vnKeyCandidates = [] } = input + const candidates = Array.isArray(vnKeyCandidates) ? vnKeyCandidates : [] + let lastError = null + + for (const vnKey of candidates) { + let vnId = 0 + let vnPhone = '' + + try { + const appointed = await appointCloudtentaclesVirtualNumber({ + ...cloudContext, + key: vnKey, + }) + vnId = Number(appointed.item?.id || 0) + vnPhone = String(appointed.item?.phone || '').trim() + + if (!vnId || !vnPhone) { + throw createHttpError('申请虚拟号成功但返回数据不完整', { + statusCode: 502, + errorCode: 'admin_task_kuaishou_cloud_invalid_vn', + }) + } + + await generateCloudtentaclesLoginCode({ + ...cloudContext, + key: vnKey, + id: vnId, + }) + + const fetchedCode = await fetchCloudtentaclesVirtualNumberCode({ + ...cloudContext, + key: vnKey, + phone: vnPhone, + }) + + await verifyCloudtentaclesLoginCode({ + ...cloudContext, + key: vnKey, + id: vnId, + code: fetchedCode.code, + }) + + const bindUrlResult = await getCloudtentaclesBindUrl({ + ...cloudContext, + key: vnKey, + id: vnId, + }) + + return { + vnKey, + vnId, + vnPhone, + bindUrl: String(bindUrlResult.bindUrl || '').trim(), + } + } catch (error) { + lastError = error + + if (vnId > 0) { + try { + await backCloudtentaclesVirtualNumber({ + ...cloudContext, + key: vnKey, + id: vnId, + }) + } catch { + // 兜底退号失败时保留原始错误,避免吞掉主链路异常 + } + } + + if (!isRecoverableKuaishouCloudVnKeyError(error)) { + throw error + } + } + } + + throw lastError || createHttpError('没有找到可用的 VN Key', { + statusCode: 409, + errorCode: 'admin_task_kuaishou_cloud_missing_binding_config', + }) +} + +function collectKuaishouCloudNameCandidates(flow) { + return Array.from(new Set([ + String(flow?.binding?.skuName || '').trim(), + String(flow?.internalSkuName || '').trim(), + String(flow?.internalSkuCode || '').trim(), + ].filter(Boolean))) +} + +function findCloudItemByNames(items, nameCandidates, preferredId = 0) { + const normalizedItems = Array.isArray(items) ? items : [] + const normalizedNames = nameCandidates + .map((item) => ({ + raw: String(item || '').trim(), + normalized: normalizeProductName(item), + })) + .filter((item) => item.raw && item.normalized) + + if (normalizedNames.length === 0 || normalizedItems.length === 0) { + return preferredId > 0 + ? normalizedItems.find((item) => Number(item.id || 0) === preferredId) || null + : null + } + + if (preferredId > 0) { + const preferred = normalizedItems.find((item) => Number(item.id || 0) === preferredId) || null + if (preferred) { + return preferred + } + } + + const exactMatches = normalizedItems.filter((item) => { + const itemName = normalizeProductName(item.name) + return normalizedNames.some((candidate) => candidate.normalized === itemName) + }) + if (exactMatches.length > 0) { + return exactMatches[0] + } + + const partialMatches = normalizedItems.filter((item) => { + const itemName = normalizeProductName(item.name) + return normalizedNames.some((candidate) => itemName.includes(candidate.normalized) || candidate.normalized.includes(itemName)) + }) + if (partialMatches.length > 0) { + return partialMatches.sort((left, right) => String(left.name || '').length - String(right.name || '').length)[0] + } + + return null +} + +function isCloudSkuLikeItem(item) { + return Boolean(item) && typeof item === 'object' && Number(item.id || 0) > 0 +} + +function isRecoverableKuaishouCloudVnKeyError(error) { + const errorCode = String(error?.errorCode || error?.code || '').trim() + const errorMessage = String(error?.message || '').trim() + return errorCode === 'cloudtentacles_vn_bind_url_failed' + && errorMessage.includes('不支持的游戏类型') +} diff --git a/apps/backend/src/services/admin/write/shared.js b/apps/backend/src/services/admin/write/shared.js new file mode 100644 index 00000000..e44ab606 --- /dev/null +++ b/apps/backend/src/services/admin/write/shared.js @@ -0,0 +1,122 @@ +// @ts-check + +import { buildClaimUrl, createTaskClaimToken } from '../../claim/claim-service.js' +import { createHttpError } from '../../../utils/http.js' +import { + canViewerConfirmAssistedRole, + canViewerRedeemAssistedTask, + isAssistedClaimTask, + parseTaskContext, +} from '../admin-read-shared-helpers.js' + +export function maskPhone(value) { + const text = String(value || '').trim() + if (!text) { + return '' + } + + if (text.length <= 7) { + return `${text.slice(0, 2)}***${text.slice(-2)}` + } + + return `${text.slice(0, 3)}****${text.slice(-4)}` +} + +export function maskCode(value) { + const text = String(value || '').trim() + if (!text) { + return '' + } + + if (text.length <= 8) { + return `${text.slice(0, 2)}***${text.slice(-2)}` + } + + return `${text.slice(0, 4)}****${text.slice(-4)}` +} + +export function resolveTaskInventoryGroupCodes(task) { + const taskContext = parseTaskContext(task) + const inventoryGroupCode = String(taskContext?.inventoryGroupCode || '').trim() + return inventoryGroupCode ? [inventoryGroupCode] : null +} + +/** + * @param {any} task + * @param {any} viewerContext + * @param {'confirm' | 'redeem'} action + */ +export function ensureViewerCanOperateAssistedTask(task, viewerContext, action) { + if (!viewerContext.canOperateAssistedTask || !isAssistedClaimTask(task)) { + throw createHttpError('当前账号没有此操作权限', { + statusCode: 403, + errorCode: 'admin_task_assisted_permission_denied', + }) + } + + if (action === 'confirm' && !canViewerConfirmAssistedRole(task, viewerContext)) { + throw createHttpError('当前任务状态还不能确认角色', { + statusCode: 409, + errorCode: 'admin_task_assisted_confirm_not_allowed', + }) + } + + if (action === 'redeem' && !canViewerRedeemAssistedTask(task, viewerContext)) { + throw createHttpError('当前任务状态还不能开始兑换', { + statusCode: 409, + errorCode: 'admin_task_assisted_redeem_not_allowed', + }) + } +} + +export function isRecoverableTaskSessionCloseError(error) { + const errorCode = String(error?.errorCode || error?.code || '').trim() + return errorCode === 'session_not_found' || errorCode === 'session_closed' +} + +export function normalizeManualDispatchOutcome(value) { + const normalized = String(value || '').trim().toLowerCase() + + if (normalized === 'failed') { + return 'failed' + } + + return 'delivered' +} + +export function getTaskClaimExpiresAt(task) { + return task?.claim_expires_at || task?.primary_claim_expires_at || null +} + +/** + * @param {any} task + */ +export async function ensureTaskClaimLink(task) { + const tokenStatus = String(task?.primary_claim_token_status || '').trim() + const token = String(task?.primary_claim_token || task?.claim_token || '').trim() + const expiredAt = getTaskClaimExpiresAt(task) + + if (tokenStatus === 'active' && token && !isClaimExpired(expiredAt)) { + return { + token, + expiredAt, + claimUrl: buildClaimUrl(token), + } + } + + const claimToken = await createTaskClaimToken(task.id) + return { + token: claimToken.token, + expiredAt: claimToken.expired_at, + claimUrl: claimToken.claimUrl, + } +} + +function isClaimExpired(expiredAt) { + if (!expiredAt) { + return false + } + + const timestamp = new Date(expiredAt).getTime() + return Number.isFinite(timestamp) && timestamp <= Date.now() +} diff --git a/apps/backend/src/services/admin/write/task-actions.js b/apps/backend/src/services/admin/write/task-actions.js new file mode 100644 index 00000000..c7d64eee --- /dev/null +++ b/apps/backend/src/services/admin/write/task-actions.js @@ -0,0 +1,514 @@ +// @ts-check + +import { + markInventoryItemDelivered, + releaseReservedInventoryItem, +} from '../../../repositories/inventory-repo.js' +import { updateClaimToken } from '../../../repositories/claim-token-repo.js' +import { listOrderItemsByOrderId } from '../../../repositories/order-item-repo.js' +import { getOrderById } from '../../../repositories/order-repo.js' +import { updateTask } from '../../../repositories/task-repo.js' +import { + getTaskInventoryBindingById, + listTaskInventoryBindingsByTaskId, +} from '../../../repositories/task-inventory-binding-repo.js' +import { createTaskEvent } from '../../../repositories/task-event-repo.js' +import { createTaskClaimToken } from '../../claim/claim-service.js' +import { confirmClaimRoleForAdminTask, redeemClaimTaskForAdminTask } from '../../claim/claim-session-service.js' +import { reserveInventoryForTask } from '../../order/inventory-service.js' +import { ensureAgisoXianyuAutoDeliveryForDeliveredTask } from '../../platforms/agiso/xianyu/auto-delivery-service.js' +import { closeTencentBrowserSession } from '../../session/session.js' +import { createHttpError } from '../../../utils/http.js' +import { nowIso } from '../../../utils/time.js' +import { + createAdminViewerContext, + getTaskPrimaryClaimTokenId, + getTaskPrimaryInventoryItemId, + isAssistedClaimTask, + isManualDispatchTask, + parseTaskContext, +} from '../admin-read-shared-helpers.js' +import { + getRequiredTask, + mapTaskActionPayload, +} from '../admin-task-read-helpers.js' +import { + ensureViewerCanOperateAssistedTask, + getTaskClaimExpiresAt, + isRecoverableTaskSessionCloseError, + maskCode, + normalizeManualDispatchOutcome, + resolveTaskInventoryGroupCodes, +} from './shared.js' + +/** @typedef {import('../../../types/admin-read-inputs.js').AdminEntityIdInput} AdminEntityIdInput */ +/** @typedef {import('../../../types/admin-read-inputs.js').AdminViewerSessionInput} AdminViewerSessionInput */ +/** @typedef {import('../../../types/admin-write-models.js').AdminTaskActionResponse} AdminTaskActionResponse */ +/** @typedef {import('../../../types/admin-write-models.js').AdminTaskBindingReleaseResponse} AdminTaskBindingReleaseResponse */ +/** @typedef {import('../../../types/admin-write-models.js').AdminTaskManualDispatchResponse} AdminTaskManualDispatchResponse */ + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +export async function releaseAdminTaskInventory(taskId) { + const task = await getRequiredTask(taskId) + const primaryInventoryItemId = getTaskPrimaryInventoryItemId(task) + + if (!primaryInventoryItemId) { + throw createHttpError('当前任务没有预占库存项', { + statusCode: 409, + errorCode: 'admin_task_no_reserved_inventory', + }) + } + + if (task.task_status === 'redeemed') { + throw createHttpError('已兑换任务不能释放库存项', { + statusCode: 409, + errorCode: 'admin_task_release_not_allowed', + }) + } + + const now = nowIso() + await releaseReservedInventoryItem(primaryInventoryItemId, now) + const updatedTask = await updateTask(task.id, { + task_status: 'waiting_inventory', + inventory_status: 'pending', + last_error: '已手动释放预占库存项', + updated_at: now, + }) + + return { + task: mapTaskActionPayload(updatedTask), + } +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +/** @param {AdminEntityIdInput} bindingId */ +export async function releaseAdminTaskInventoryBinding(taskId, bindingId) { + const task = await getRequiredTask(taskId) + const binding = await getTaskInventoryBindingById(Number(bindingId)) + + if (!binding || Number(binding.task_id) !== Number(task.id)) { + throw createHttpError('任务库存绑定不存在', { + statusCode: 404, + errorCode: 'admin_task_inventory_binding_not_found', + }) + } + + if (String(binding.binding_status || '').trim() !== 'reserved') { + throw createHttpError('当前库存绑定不是预占状态,不能释放', { + statusCode: 409, + errorCode: 'admin_task_inventory_binding_release_not_allowed', + }) + } + + if (task.task_status === 'redeemed') { + throw createHttpError('已兑换任务不能释放库存绑定', { + statusCode: 409, + errorCode: 'admin_task_release_not_allowed', + }) + } + + const now = nowIso() + await releaseReservedInventoryItem(binding.inventory_item_id, now) + + const remainingBindings = await listTaskInventoryBindingsByTaskId(task.id) + const activeBindings = remainingBindings.filter((item) => ['reserved', 'consumed'].includes(String(item.binding_status || '').trim())) + const hasReservedBindings = activeBindings.some((item) => String(item.binding_status || '').trim() === 'reserved') + const hasConsumedBindings = activeBindings.some((item) => String(item.binding_status || '').trim() === 'consumed') + const nextInventoryStatus = hasReservedBindings ? 'reserved' : (hasConsumedBindings ? 'consumed' : 'pending') + const nextTaskStatus = !activeBindings.length && !['closed', 'expired'].includes(String(task.task_status || '').trim()) + ? 'waiting_inventory' + : task.task_status + + if (!hasReservedBindings) { + const primaryClaimTokenId = getTaskPrimaryClaimTokenId(task) + if (primaryClaimTokenId) { + await updateClaimToken(primaryClaimTokenId, { + status: 'revoked', + updated_at: now, + }) + } + } + + const updatedTask = await updateTask(task.id, { + task_status: nextTaskStatus, + inventory_status: nextInventoryStatus, + last_error: !activeBindings.length ? '已手动释放预占库存绑定' : (task.last_error || ''), + updated_at: now, + }) + + await createTaskEvent(task.id, 'inventory_binding_released', { + bindingId: Number(binding.id), + inventoryItemId: Number(binding.inventory_item_id), + roleKey: String(binding.role_key || '').trim(), + }, now) + + return { + task: mapTaskActionPayload(updatedTask), + bindingId: Number(binding.id), + inventoryItemId: Number(binding.inventory_item_id), + } +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +/** @param {AdminViewerSessionInput | null} [session] */ +export async function regenerateAdminTaskClaimLink(taskId, session = null) { + const task = await getRequiredTask(taskId) + const primaryClaimTokenId = getTaskPrimaryClaimTokenId(task) + const viewerContext = createAdminViewerContext(session) + + if (isManualDispatchTask(task)) { + throw createHttpError('人工履约任务不需要领取链接,请直接回写人工履约结果', { + statusCode: 409, + errorCode: 'admin_task_manual_dispatch_claim_not_allowed', + }) + } + + if (['redeemed', 'closed'].includes(task.task_status)) { + throw createHttpError('当前任务状态不允许重新生成领取链接', { + statusCode: 409, + errorCode: 'admin_task_regenerate_not_allowed', + }) + } + + if (!viewerContext.canManageTaskLifecycle && !isAssistedClaimTask(task)) { + throw createHttpError('当前账号只能重发半自动客服任务的领取链接', { + statusCode: 403, + errorCode: 'admin_task_regenerate_permission_denied', + }) + } + + const now = nowIso() + if (primaryClaimTokenId) { + await updateClaimToken(primaryClaimTokenId, { + status: 'revoked', + updated_at: now, + }) + } + + const claimToken = await createTaskClaimToken(task.id) + const updatedTask = await updateTask(task.id, { + claim_token: claimToken.token, + claim_expires_at: claimToken.expired_at, + task_status: 'link_generated', + last_error: '', + updated_at: now, + }) + + return { + task: mapTaskActionPayload(updatedTask), + claimUrl: claimToken.claimUrl, + token: claimToken.token, + } +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +/** @param {AdminViewerSessionInput | null} [session] */ +export async function confirmAdminTaskAssistedRole(taskId, session = null) { + const task = await getRequiredTask(taskId) + const viewerContext = createAdminViewerContext(session) + ensureViewerCanOperateAssistedTask(task, viewerContext, 'confirm') + await confirmClaimRoleForAdminTask(task.id) + const updatedTask = await getRequiredTask(task.id) + + return { + task: mapTaskActionPayload(updatedTask), + } +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +/** @param {AdminViewerSessionInput | null} [session] */ +export async function redeemAdminTaskAssisted(taskId, session = null) { + const task = await getRequiredTask(taskId) + const viewerContext = createAdminViewerContext(session) + ensureViewerCanOperateAssistedTask(task, viewerContext, 'redeem') + await redeemClaimTaskForAdminTask(task.id) + const updatedTask = await getRequiredTask(task.id) + + return { + task: mapTaskActionPayload(updatedTask), + } +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +export async function closeAdminTask(taskId) { + return closeAdminTaskWithDeps(taskId) +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +export async function closeAdminTaskWithDeps( + taskId, + { + getRequiredTask: getTask = getRequiredTask, + listTaskInventoryBindingsByTaskId: listBindings = listTaskInventoryBindingsByTaskId, + releaseReservedInventoryItem: releaseReserved = releaseReservedInventoryItem, + updateClaimToken: updateToken = updateClaimToken, + updateTask: updateTaskRecord = updateTask, + createTaskEvent: createEvent = createTaskEvent, + closeTencentBrowserSession: closeSession = closeTencentBrowserSession, + nowIso: getNowIso = nowIso, + } = {}, +) { + const task = await getTask(taskId) + + if (task.task_status === 'redeemed') { + throw createHttpError('已兑换任务不能关闭', { + statusCode: 409, + errorCode: 'admin_task_close_not_allowed', + }) + } + + const now = getNowIso() + const bindings = await listBindings(task.id) + const reservedBindings = bindings.filter((binding) => String(binding.binding_status || '').trim() === 'reserved') + const releasedInventoryItemIds = Array.from(new Set(reservedBindings.map((binding) => Number(binding.inventory_item_id)).filter((id) => id > 0))) + const hasConsumedBindings = bindings.some((binding) => String(binding.binding_status || '').trim() === 'consumed') + const primaryClaimTokenId = getTaskPrimaryClaimTokenId(task) + let browserSessionClosed = false + + if (primaryClaimTokenId) { + await updateToken(primaryClaimTokenId, { + status: 'revoked', + expired_at: now, + updated_at: now, + }) + } + + if (task.browser_session_id) { + try { + await closeSession(task.browser_session_id) + browserSessionClosed = true + } catch (error) { + if (!isRecoverableTaskSessionCloseError(error)) { + throw error + } + } + } + + for (const inventoryItemId of releasedInventoryItemIds) { + await releaseReserved(inventoryItemId, now) + } + + const closeReasonParts = ['已手动关闭任务'] + if (primaryClaimTokenId) { + closeReasonParts.push('领取链接已失效') + } + if (releasedInventoryItemIds.length > 0) { + closeReasonParts.push('预占库存已释放') + } + + const updatedTask = await updateTaskRecord(task.id, { + task_status: 'closed', + inventory_status: hasConsumedBindings ? 'consumed' : 'pending', + delivery_status: 'closed', + user_action_status: 'closed', + claim_expires_at: primaryClaimTokenId ? now : task.claim_expires_at, + browser_session_id: '', + last_error: task.last_error || closeReasonParts.join(','), + updated_at: now, + }) + + await createEvent(task.id, 'task_closed', { + claimTokenRevoked: Boolean(primaryClaimTokenId), + releasedInventoryItemIds, + releasedInventoryCount: releasedInventoryItemIds.length, + browserSessionClosed, + }, now) + + return { + task: mapTaskActionPayload(updatedTask), + } +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +export async function markAdminTaskManualReview(taskId) { + const task = await getRequiredTask(taskId) + const updatedTask = await updateTask(task.id, { + task_status: 'manual_review', + delivery_status: task.delivery_status || 'pending', + last_error: task.last_error || '已转人工处理', + updated_at: nowIso(), + }) + + return { + task: mapTaskActionPayload(updatedTask), + } +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +export async function retryAdminTask(taskId) { + const task = await getRequiredTask(taskId) + const now = nowIso() + + if (isManualDispatchTask(task)) { + throw createHttpError('人工履约任务不能走自动重试,请在详情页直接回写人工履约结果', { + statusCode: 409, + errorCode: 'admin_task_manual_dispatch_retry_not_allowed', + }) + } + + if (!['retry_pending', 'manual_review', 'waiting_inventory'].includes(task.task_status)) { + throw createHttpError('当前任务状态不允许重试', { + statusCode: 409, + errorCode: 'admin_task_retry_not_allowed', + }) + } + + const taskContext = parseTaskContext(task) + const primaryRequirement = taskContext.primaryRequirement || null + const orderItems = await listOrderItemsByOrderId(task.order_id) + const orderItem = orderItems.find((item) => item.id === task.order_item_id) || null + let reservedInventoryItemId = getTaskPrimaryInventoryItemId(task) + let claimTokenId = getTaskPrimaryClaimTokenId(task) + let nextStatus = 'link_generated' + let lastError = '' + let claimExpiresAt = getTaskClaimExpiresAt(task) + let claimUrl = '' + let token = '' + + if (!reservedInventoryItemId) { + const reserved = await reserveInventoryForTask({ + skuCode: orderItem?.sku_code || '', + taskId: task.id, + credentialType: primaryRequirement?.credentialType || 'tencent_code', + roleKey: primaryRequirement?.roleKey || 'primary_code', + inventoryGroupCodes: resolveTaskInventoryGroupCodes(task), + }) + + if (!reserved) { + nextStatus = 'waiting_inventory' + lastError = '库存不足,等待可用库存凭据' + } else { + reservedInventoryItemId = reserved.id + } + } + + if (nextStatus === 'link_generated' && !claimTokenId) { + const claimToken = await createTaskClaimToken(task.id) + claimTokenId = claimToken.id + claimExpiresAt = claimToken.expired_at + claimUrl = claimToken.claimUrl + token = claimToken.token + } + + const updatedTask = await updateTask(task.id, { + task_status: nextStatus, + inventory_status: reservedInventoryItemId ? 'reserved' : 'pending', + claim_token: token || task.claim_token || '', + claim_expires_at: claimExpiresAt, + last_error: lastError, + updated_at: now, + }) + + const response = { + task: mapTaskActionPayload(updatedTask), + } + + if (claimUrl) { + response.claimUrl = claimUrl + response.token = token + } + + return response +} + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} taskId */ +/** @param {{ outcome?: string, resultMessage?: string, deliveryReference?: string, deliveredCredential?: string }} [payload] */ +/** @param {AdminViewerSessionInput | null} [session] */ +export async function completeAdminTaskManualDispatch(taskId, payload = {}, session = null) { + const task = await getRequiredTask(taskId) + + if (!isManualDispatchTask(task)) { + throw createHttpError('当前任务不是人工履约任务', { + statusCode: 409, + errorCode: 'admin_task_not_manual_dispatch', + }) + } + + if (['redeemed', 'closed'].includes(task.task_status)) { + throw createHttpError('当前任务已经完结,不能重复回写人工履约结果', { + statusCode: 409, + errorCode: 'admin_task_manual_dispatch_already_completed', + }) + } + + const now = nowIso() + const outcome = normalizeManualDispatchOutcome(payload.outcome) + const resultMessage = String(payload.resultMessage || '').trim() + const deliveryReference = String(payload.deliveryReference || '').trim() + const deliveredCredential = String(payload.deliveredCredential || '').trim() + const context = parseTaskContext(task) + const inventoryItemId = getTaskPrimaryInventoryItemId(task) + const resultCode = outcome === 'failed' ? 'manual_dispatch_failed' : 'manual_dispatch_delivered' + const fallbackMessage = outcome === 'failed' ? '人工履约失败' : '人工履约已完成' + const nextTaskStatus = outcome === 'failed' ? 'closed' : 'redeemed' + const nextDeliveryStatus = outcome === 'failed' ? 'failed' : 'delivered' + const nextInventoryStatus = outcome === 'delivered' && inventoryItemId + ? 'consumed' + : task.inventory_status || 'not_required' + const manualDispatch = { + outcome, + deliveryReference, + deliveredCredential, + resultMessage: resultMessage || fallbackMessage, + completedAt: now, + completedBy: session + ? { + userId: Number(session.userId || 0), + username: String(session.username || ''), + role: String(session.role || ''), + } + : null, + } + + if (outcome === 'delivered' && inventoryItemId) { + await markInventoryItemDelivered(inventoryItemId, now) + } + + const updatedTask = await updateTask(task.id, { + task_status: nextTaskStatus, + inventory_status: nextInventoryStatus, + delivery_status: nextDeliveryStatus, + result_code: resultCode, + result_message: resultMessage || fallbackMessage, + user_action_status: 'not_required', + last_error: outcome === 'failed' ? (resultMessage || fallbackMessage) : '', + redeemed_at: outcome === 'delivered' ? now : task.redeemed_at || null, + context_json: JSON.stringify({ + ...context, + manualDispatch, + }), + updated_at: now, + }) + + await createTaskEvent(task.id, 'manual_dispatch_completed', { + outcome, + resultCode, + resultMessage: resultMessage || fallbackMessage, + deliveryReference, + deliveredCredentialMasked: maskCode(deliveredCredential), + completedBy: manualDispatch.completedBy, + }, now) + + const taskAfterAutoDelivery = outcome === 'delivered' + ? (await ensureAgisoXianyuAutoDeliveryForDeliveredTask({ + order: await getOrderById(task.order_id), + task: updatedTask, + trigger: 'manual_dispatch_completed', + })).task || updatedTask + : updatedTask + + return { + outcome, + task: mapTaskActionPayload(taskAfterAutoDelivery), + } +} diff --git a/apps/backend/src/services/admin/write/webhook-events.js b/apps/backend/src/services/admin/write/webhook-events.js new file mode 100644 index 00000000..72acb1e0 --- /dev/null +++ b/apps/backend/src/services/admin/write/webhook-events.js @@ -0,0 +1,36 @@ +// @ts-check + +import { getWebhookEventById } from '../../../repositories/webhook-event-repo.js' +import { replayAgisoTradeWebhookEvent } from '../../order/webhook-service.js' +import { createHttpError } from '../../../utils/http.js' + +/** @typedef {import('../../../types/admin-read-inputs.js').AdminEntityIdInput} AdminEntityIdInput */ +/** @typedef {import('../../../types/admin-write-models.js').AdminWebhookReplayResponse} AdminWebhookReplayResponse */ + +/** @returns {Promise} */ +/** @param {AdminEntityIdInput} eventId */ +export async function replayAdminWebhookEvent(eventId) { + const event = await getWebhookEventById(Number(eventId)) + + if (!event) { + throw createHttpError('Webhook 事件不存在', { + statusCode: 404, + errorCode: 'admin_webhook_not_found', + }) + } + + if (String(event.provider || event.platform || '').trim() !== 'agiso') { + throw createHttpError('当前只支持重放 agiso webhook', { + statusCode: 409, + errorCode: 'admin_webhook_replay_not_supported', + }) + } + + const result = await replayAgisoTradeWebhookEvent(event) + + return { + eventId: event.id, + replayed: true, + result, + } +}