// @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 { 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 { appointCloudtentaclesVirtualNumber, backCloudtentaclesVirtualNumber, 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 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, }, 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', 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), } } /** @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 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 || flow.ticket.code, 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), 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 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 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, }, }, } const updatedTask = await updateTask(task.id, { task_status: 'completed', delivery_status: 'delivered', result_code: 'kuaishou_cloud_completed', result_message: 'cloudtentacles 发货并退号完成', redeemed_at: now, last_error: '', 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) 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 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(), capturedAt: ticket.capturedAt || null, capturedBy: ticket.capturedBy || null, }, 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, }, 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(), autoConsumeEnabled: consume.autoConsumeEnabled === true, }, } } 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 }