// @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 { replayAgisoTradeWebhookEvent } from '../order/webhook-service.js' import { ensureAgisoXianyuAutoDeliveryForDeliveredTask } from '../platforms/agiso/xianyu/auto-delivery-service.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-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 */ /** @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' 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, createdAt: now, updatedAt: now, }, ]) if (created === 0) { throw createHttpError('库存凭据已存在,不能重复新增', { statusCode: 409, errorCode: 'admin_inventory_duplicate', }) } const { items } = await listInventoryItems({ page: 1, pageSize: 1, skuCode, }) const createdItem = items.find((item) => item.display_value === displayValue) || 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, 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) { const task = await getRequiredTask(taskId) if (task.task_status === 'redeemed') { throw createHttpError('已兑换任务不能关闭', { statusCode: 409, errorCode: 'admin_task_close_not_allowed', }) } const updatedTask = await updateTask(task.id, { task_status: 'closed', delivery_status: 'closed', last_error: task.last_error || '已手动关闭任务', updated_at: nowIso(), }) 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', }) 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), } } /** * @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' if (!skuCode || !displayValue) { continue } rows.push({ skuCode, displayValue, batchNo, credentialType }) } } 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' 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 }) } } 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 } /** * @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 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 } 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)}` }