删除咸鱼旧链路

This commit is contained in:
yml
2026-05-25 21:46:39 +08:00
parent 9aafddd624
commit e7aa194dde
142 changed files with 539 additions and 17373 deletions
@@ -1,273 +0,0 @@
import test from 'node:test'
import assert from 'node:assert/strict'
import { redeemClaimTaskWithInventoryFallbackWithDeps } from './claim-session-service.js'
test('redeemClaimTaskWithInventoryFallbackWithDeps invalidates missing CDKEY and retries with replacement inventory', async () => {
const events = []
const invalidated = []
const reserveCalls = []
const reloadCalls = []
const context = {
task: {
id: 42,
browser_session_id: 'txbs-test',
},
orderItem: {
sku_code: 'df-cdk',
},
}
const initialInventoryItem = {
id: 1,
display_value: 'DJQFf7HgONCL4XD4EH',
credential_type: 'tencent_code',
inventory_group_code: 'A组',
}
const replacementInventoryItem = {
id: 2,
display_value: 'ABCD1234EFGH5678',
credential_type: 'tencent_code',
inventory_group_code: 'A组',
}
let redeemCount = 0
const result = await redeemClaimTaskWithInventoryFallbackWithDeps(
context,
initialInventoryItem,
{
reloadTencentBrowserSession: async (sessionId) => {
reloadCalls.push(sessionId)
return { sessionId, status: 'ready_to_redeem' }
},
redeemTencentBrowserSession: async (sessionId, payload) => {
redeemCount += 1
if (redeemCount === 1) {
assert.equal(sessionId, 'txbs-test')
assert.equal(payload.code, 'DJQFf7HgONCL4XD4EH')
return {
redeem: {
final: {
redeem: {
iRet: -183,
sMsg: '该CDKEY不存在,请您确认后输入!',
},
},
},
}
}
assert.equal(payload.code, 'ABCD1234EFGH5678')
return {
redeem: {
final: {
redeem: {
iRet: 0,
sMsg: '领取成功',
},
},
},
}
},
invalidateReservedInventoryItem: async (inventoryItemId, reason, updatedAt) => {
invalidated.push({ inventoryItemId, reason, updatedAt })
return { id: inventoryItemId, status: 'invalid', invalid_reason: reason, updated_at: updatedAt }
},
reserveInventoryForTask: async (payload) => {
reserveCalls.push(payload)
return replacementInventoryItem
},
createTaskEvent: async (taskId, eventType, payload, createdAt) => {
events.push({ taskId, eventType, payload, createdAt })
},
nowIso: (() => {
const values = [
'2026-04-14T12:03:34.000Z',
'2026-04-14T12:03:35.000Z',
]
let index = 0
return () => values[index++] || values[values.length - 1]
})(),
},
)
assert.equal(result.inventoryItem.id, 2)
assert.equal(result.classification.outcome, 'success')
assert.equal(result.attempts.length, 2)
assert.deepEqual(
result.attempts.map((attempt) => ({
attempt: attempt.attempt,
inventoryItemId: attempt.inventoryItemId,
outcome: attempt.outcome,
resultCode: attempt.resultCode,
})),
[
{
attempt: 1,
inventoryItemId: 1,
outcome: 'code_invalid',
resultCode: '-183',
},
{
attempt: 2,
inventoryItemId: 2,
outcome: 'success',
resultCode: '0',
},
],
)
assert.deepEqual(invalidated, [
{
inventoryItemId: 1,
reason: '该CDKEY不存在,请您确认后输入!',
updatedAt: '2026-04-14T12:03:34.000Z',
},
])
assert.deepEqual(reserveCalls, [
{
skuCode: 'df-cdk',
taskId: 42,
credentialType: 'tencent_code',
roleKey: 'primary_code',
inventoryGroupCodes: ['A组'],
},
])
assert.deepEqual(reloadCalls, ['txbs-test'])
assert.deepEqual(
events.map((event) => ({
eventType: event.eventType,
taskId: event.taskId,
createdAt: event.createdAt,
payload: event.payload,
})),
[
{
eventType: 'claim_redeem_code_invalid',
taskId: 42,
createdAt: '2026-04-14T12:03:34.000Z',
payload: {
inventoryItemId: 1,
codeMasked: 'DJQF****D4EH',
resultCode: '-183',
resultMessage: '该CDKEY不存在,请您确认后输入!',
},
},
{
eventType: 'claim_redeem_inventory_replaced',
taskId: 42,
createdAt: '2026-04-14T12:03:35.000Z',
payload: {
previousInventoryItemId: 1,
previousCodeMasked: 'DJQF****D4EH',
previousOutcome: 'code_invalid',
nextInventoryItemId: 2,
nextCodeMasked: 'ABCD****5678',
credentialType: 'tencent_code',
inventoryGroupCode: 'A组',
},
},
],
)
})
test('redeemClaimTaskWithInventoryFallbackWithDeps waits for inventory when used code has no replacement', async () => {
const events = []
const consumed = []
const reserveCalls = []
const context = {
task: {
id: 57,
browser_session_id: 'txbs-used',
},
orderItem: {
sku_code: 'df-cdk-used',
},
}
const initialInventoryItem = {
id: 7,
display_value: 'USED1234CODE5678',
credential_type: 'tencent_code',
inventory_group_code: 'B组',
}
let thrown = null
try {
await redeemClaimTaskWithInventoryFallbackWithDeps(
context,
initialInventoryItem,
{
redeemTencentBrowserSession: async () => ({
redeem: {
final: {
redeem: {
popup: {
text: '兑换结果',
detail: '兑换码已使用。',
},
},
},
},
}),
markInventoryItemConsumed: async (inventoryItemId, reason, updatedAt) => {
consumed.push({ inventoryItemId, reason, updatedAt })
return { id: inventoryItemId, status: 'consumed', invalid_reason: reason, updated_at: updatedAt }
},
reserveInventoryForTask: async (payload) => {
reserveCalls.push(payload)
return null
},
createTaskEvent: async (taskId, eventType, payload, createdAt) => {
events.push({ taskId, eventType, payload, createdAt })
},
nowIso: () => '2026-04-14T12:05:00.000Z',
},
)
} catch (error) {
thrown = error
}
assert.ok(thrown)
assert.equal(thrown.errorCode, 'claim_inventory_replacement_exhausted')
assert.equal(thrown.redeemTaskState.taskStatus, 'waiting_inventory')
assert.equal(thrown.redeemTaskState.inventoryStatus, 'pending')
assert.equal(thrown.redeemTaskState.deliveryStatus, 'pending')
assert.equal(thrown.redeemTaskState.attempts.length, 1)
assert.deepEqual(consumed, [
{
inventoryItemId: 7,
reason: '兑换码已使用。',
updatedAt: '2026-04-14T12:05:00.000Z',
},
])
assert.deepEqual(reserveCalls, [
{
skuCode: 'df-cdk-used',
taskId: 57,
credentialType: 'tencent_code',
roleKey: 'primary_code',
inventoryGroupCodes: ['B组'],
},
])
assert.deepEqual(
events.map((event) => ({
taskId: event.taskId,
eventType: event.eventType,
payload: event.payload,
createdAt: event.createdAt,
})),
[
{
taskId: 57,
eventType: 'claim_redeem_code_used',
payload: {
inventoryItemId: 7,
codeMasked: 'USED****5678',
resultCode: '',
resultMessage: '兑换码已使用。',
},
createdAt: '2026-04-14T12:05:00.000Z',
},
],
)
})
@@ -1,28 +0,0 @@
export {
getClaimDetail,
getClaimDetailForAdminTask,
getClaimSessionSummary,
getClaimSessionSummaryForAdminTask,
confirmClaimRole,
confirmClaimRoleForAdminTask,
redeemClaimTask,
redeemClaimTaskForAdminTask,
getClaimScreenshotPath,
} from './session/actions.js'
export {
createClaimSession,
createClaimSessionForAdminTask,
reloadClaimSession,
reloadClaimSessionForAdminTask,
closeClaimSession,
closeClaimSessionForAdminTask,
} from './session/lifecycle.js'
export {
getClaimContext,
} from './session/context.js'
export {
redeemClaimTaskWithInventoryFallbackWithDeps,
} from './session/redeem.js'
@@ -1,32 +1,95 @@
import { buildClaimUrl } from '../claim-service.js'
import { createHttpError } from '../../../utils/http.js'
import { maskCode as maskCodeValue } from '../../../utils/masking.js'
import { formatFenToAmount, normalizeFen } from '../../../utils/money.js'
import {
parseTaskContext as parseTaskContextValue,
parseTaskState as parseTaskStateValue,
} from '../../../utils/task-json.js'
import { findClaimTokenByToken, updateClaimToken } from '../../repositories/claim-token-repo.js'
import { releaseReservedInventoryItem } from '../../repositories/inventory-repo.js'
import { getOrderItemById } from '../../repositories/order-item-repo.js'
import { getOrderById } from '../../repositories/order-repo.js'
import { findTaskByClaimTokenId, updateTask } from '../../repositories/task-repo.js'
import { createHttpError } from '../../utils/http.js'
import { formatFenToAmount, normalizeFen } from '../../utils/money.js'
import { parseTaskContext as parseTaskContextValue } from '../../utils/task-json.js'
import { nowIso } from '../../utils/time.js'
import { buildClaimUrl } from './claim-service.js'
export const CLAIM_TERMINAL_STATUSES = new Set(['expired', 'closed'])
export const REDEEM_REPLACEMENT_LIMIT = 10
export const KUAISHOU_CLOUD_CLAIM_GUIDE_BASE_PATH = '/kuaishou-cloud-guide'
export function buildClaimDetailPayload({ claimToken, task, order, orderItem, session }) {
const screenshotReady = Boolean(task.screenshot_path) || Boolean(session?.artifacts?.hasScreenshot)
const screenshotUrl = screenshotReady ? `/api/v1/claim/${claimToken.token}/screenshot` : ''
const finalRedeem = session?.redeem?.final?.redeem || null
export async function getClaimContext(token: unknown) {
const normalized = String(token || '').trim()
if (!normalized) {
throw createHttpError('缺少领取 token', {
statusCode: 400,
errorCode: 'missing_claim_token',
})
}
const claimToken = await findClaimTokenByToken(normalized)
if (!claimToken) {
throw createHttpError('领取链接无效或不存在', {
statusCode: 404,
errorCode: 'claim_token_not_found',
})
}
const task = await findTaskByClaimTokenId(claimToken.id)
if (!task) {
throw createHttpError('领取任务不存在', {
statusCode: 404,
errorCode: 'claim_task_not_found',
})
}
if (claimToken.status !== 'active') {
throw createHttpError('领取链接当前不可用', {
statusCode: 410,
errorCode: 'claim_token_inactive',
})
}
if (claimToken.expired_at && new Date(claimToken.expired_at).getTime() <= Date.now()) {
const expiredContext = await expireClaimContext(claimToken, task)
throw createHttpError('领取链接已过期', {
statusCode: 410,
errorCode: 'claim_token_expired',
context: expiredContext,
})
}
const [order, orderItem] = await Promise.all([
getOrderById(task.order_id),
getOrderItemById(task.order_item_id),
])
if (!order || !orderItem) {
throw createHttpError('领取任务关联订单不完整', {
statusCode: 500,
errorCode: 'claim_order_incomplete',
})
}
return {
claimToken,
task,
order,
orderItem,
}
}
export function buildClaimDetailPayload({ claimToken, task, order, orderItem }) {
const kuaishouCloudFulfillment = mapClaimKuaishouCloudFulfillment(task, order)
return {
tokenStatus: claimToken.status,
claimUrl: buildClaimUrl(claimToken.token),
flowType: kuaishouCloudFulfillment ? 'kuaishou_cloud' : 'tencent_claim',
flowType: 'kuaishou_cloud',
task: {
taskId: task.id,
taskNo: task.task_no,
status: task.task_status,
executorKey: task.executor_key || '',
requiresSupportReview: isAssistedClaimTask(task),
requiresSupportReview: false,
expiresAt: claimToken.expired_at,
claimedAt: task.claimed_at,
roleConfirmedAt: task.role_confirmed_at,
@@ -51,14 +114,14 @@ export function buildClaimDetailPayload({ claimToken, task, order, orderItem, se
skuName: orderItem.sku_name,
quantity: orderItem.quantity,
},
session,
session: null,
kuaishouCloudFulfillment,
result: task.redeemed_at || session?.status === 'redeemed'
result: task.redeemed_at
? {
resultCode: String(task.result_code || finalRedeem?.iRet || finalRedeem?.ret || ''),
resultMessage: String(task.result_message || finalRedeem?.sMsg || finalRedeem?.msg || session?.notice || ''),
screenshotReady,
screenshotUrl,
resultCode: String(task.result_code || ''),
resultMessage: String(task.result_message || ''),
screenshotReady: false,
screenshotUrl: '',
}
: null,
}
@@ -151,58 +214,35 @@ export function mapClaimKuaishouCloudFulfillment(task, order) {
}
}
export function assertTaskCanProceed(task) {
if (CLAIM_TERMINAL_STATUSES.has(String(task.task_status || ''))) {
throw createHttpError('当前任务已经结束,不能继续操作', {
statusCode: 410,
errorCode: 'claim_task_closed',
})
}
}
export function assertPublicClaimActionAllowed(task, action) {
if (!isAssistedClaimTask(task)) {
return
}
const message = action === 'confirm'
? '当前商品需要客服复核角色,请登录后联系人工继续'
: '当前商品需要客服确认后再执行兑换,请联系人工继续'
throw createHttpError(message, {
statusCode: 409,
errorCode: 'claim_support_review_required',
})
}
export function mergeTaskContext(task, patch = {}) {
return {
...parseTaskContext(task),
...patch,
}
}
export function parseTaskState(task) {
return parseTaskStateValue(task)
}
export function isAssistedClaimTask(task) {
return String(task?.executor_key || '').trim() === 'tencent_claim_assisted'
}
export function parseTaskContext(task) {
function parseTaskContext(task) {
return parseTaskContextValue(task)
}
export function normalizeClaimLoginType(loginType) {
return String(loginType || '').trim() === 'wx' ? 'wx' : 'qq'
}
async function expireClaimContext(claimToken, task) {
const now = nowIso()
const nextClaimToken = await updateClaimToken(claimToken.id, {
status: 'expired',
updated_at: now,
})
export function isRecoverableSessionError(error) {
const errorCode = String(error?.errorCode || error?.code || '').trim()
return errorCode === 'session_not_found' || errorCode === 'session_closed'
}
let nextTask = task
export function maskCode(value) {
return maskCodeValue(value, { shortMask: '****' })
if (!CLAIM_TERMINAL_STATUSES.has(String(task.task_status || '')) && task.task_status !== 'redeemed') {
if (task.primary_inventory_item_id) {
await releaseReservedInventoryItem(task.primary_inventory_item_id, now)
}
nextTask = await updateTask(task.id, {
task_status: 'expired',
inventory_status: 'pending',
user_action_status: 'expired',
last_error: '领取链接已过期,预占库存项已释放',
updated_at: now,
})
}
return {
claimToken: nextClaimToken,
task: nextTask,
}
}
@@ -18,8 +18,7 @@ import {
prepareKuaishouCloudFulfillmentTask,
refreshKuaishouCloudTaskRoleInfo,
} from '../fulfillment/kuaishou-cloud-task-service.js'
import { buildClaimDetailPayload } from './session/shared.js'
import { getClaimContext } from './session/context.js'
import { buildClaimDetailPayload, getClaimContext } from './kuaishou-cloud-claim-context.js'
import { syncKuaishouCloudRoleInfo } from './kuaishou-cloud-sync-service.js'
const KUAISHOU_CLOUD_GUIDE_DIR = path.resolve(PROJECT_ROOT, '../../tems/imgs')
@@ -170,7 +169,6 @@ export async function getKuaishouCloudClaimDetail(token: unknown) {
task,
order: context.order,
orderItem: context.orderItem,
session: null,
})
}
@@ -1,312 +0,0 @@
import {
getTencentBrowserSession,
getTencentBrowserSessionScreenshotPath,
} from '../../session/session.js'
import { createTaskEvent } from '../../../repositories/task-event-repo.js'
import { updateTask } from '../../../repositories/task-repo.js'
import {
getInventoryItemById,
markInventoryItemDelivered,
} from '../../../repositories/inventory-repo.js'
import { ensureAgisoXianyuAutoDeliveryForDeliveredTask } from '../../platforms/agiso/xianyu/auto-delivery-service.js'
import { syncKuaishouCloudRoleInfo } from '../kuaishou-cloud-sync-service.js'
import { notifyClaimRedeemNeedsAttention } from '../../notification/domain-notifications.js'
import { createHttpError } from '../../../utils/http.js'
import { nowIso } from '../../../utils/time.js'
import { getClaimContext, getClaimContextByTaskId } from './context.js'
import { loadTaskSession, syncTaskWithSession } from './runtime.js'
import { redeemClaimTaskWithInventoryFallback, resolveTencentRedeemResultCode } from './redeem.js'
import {
assertPublicClaimActionAllowed,
assertTaskCanProceed,
buildClaimDetailPayload,
maskCode,
mergeTaskContext,
} from './shared.js'
export async function getClaimDetail(token, { includeQrImage = true } = {}) {
const context = await getClaimContext(token)
let { task, session } = await loadTaskSession(context.task, { includeQrImage })
if (String(task.executor_key || '').trim() === 'kuaishou_ct_assisted') {
task = await syncKuaishouCloudRoleInfo(task)
}
const syncedTask = session ? await syncTaskWithSession(task, session) : task
return buildClaimDetailPayload({
claimToken: context.claimToken,
task: syncedTask,
order: context.order,
orderItem: context.orderItem,
session,
})
}
export async function getClaimDetailForAdminTask(taskId, { includeQrImage = true } = {}) {
const context = await getClaimContextByTaskId(taskId)
const { task, session } = await loadTaskSession(context.task, { includeQrImage })
const syncedTask = session ? await syncTaskWithSession(task, session) : task
return buildClaimDetailPayload({
claimToken: context.claimToken,
task: syncedTask,
order: context.order,
orderItem: context.orderItem,
session,
})
}
export async function getClaimSessionSummary(token) {
return getClaimSessionSummaryByContextLoader(() => getClaimContext(token))
}
export async function getClaimSessionSummaryForAdminTask(taskId) {
return getClaimSessionSummaryByContextLoader(() => getClaimContextByTaskId(taskId))
}
export async function confirmClaimRole(token) {
const context = await getClaimContext(token)
assertPublicClaimActionAllowed(context.task, 'confirm')
return finalizeClaimRoleConfirmation(context)
}
export async function confirmClaimRoleForAdminTask(taskId) {
const context = await getClaimContextByTaskId(taskId)
return finalizeClaimRoleConfirmation(context)
}
export async function redeemClaimTask(token) {
const context = await getClaimContext(token)
assertTaskCanProceed(context.task)
assertPublicClaimActionAllowed(context.task, 'redeem')
return finalizeClaimTaskRedeem(context)
}
export async function redeemClaimTaskForAdminTask(taskId) {
const context = await getClaimContextByTaskId(taskId)
assertTaskCanProceed(context.task)
return finalizeClaimTaskRedeem(context)
}
export async function getClaimScreenshotPath(token) {
const context = await getClaimContext(token)
if (context.task.screenshot_path) {
return context.task.screenshot_path
}
if (!context.task.browser_session_id) {
throw createHttpError('当前任务还没有兑换截图', {
statusCode: 404,
errorCode: 'claim_screenshot_not_ready',
})
}
return getTencentBrowserSessionScreenshotPath(context.task.browser_session_id)
}
async function getClaimSessionSummaryByContextLoader(loadContext) {
const context = await loadContext()
const { task, session } = await loadTaskSession(context.task, {
includeQrImage: false,
})
if (!session) {
return buildClaimDetailPayload({
claimToken: context.claimToken,
task,
order: context.order,
orderItem: context.orderItem,
session: null,
})
}
const syncedTask = await syncTaskWithSession(task, session)
return buildClaimDetailPayload({
claimToken: context.claimToken,
task: syncedTask,
order: context.order,
orderItem: context.orderItem,
session,
})
}
async function finalizeClaimRoleConfirmation(context) {
if (!context.task.browser_session_id) {
throw createHttpError('当前任务还没有创建浏览器会话', {
statusCode: 409,
errorCode: 'claim_session_not_created',
})
}
const session = await getTencentBrowserSession(context.task.browser_session_id)
const activityInfo = session.activityInfo || null
if (!activityInfo?.role?.ready) {
throw createHttpError('当前角色信息还没有准备好', {
statusCode: 409,
errorCode: 'claim_role_not_ready',
})
}
const updatedTask = await updateTask(context.task.id, {
task_status: 'role_confirmed',
nickname: String(activityInfo.nickname || ''),
role_id: String(activityInfo.role.roleId || ''),
role_name: String(activityInfo.role.roleName || ''),
area: String(activityInfo.role.area || ''),
partition_name: String(activityInfo.role.partition || ''),
user_action_status: 'role_confirmed',
role_confirmed_at: nowIso(),
updated_at: nowIso(),
last_error: '',
})
return buildClaimDetailPayload({
claimToken: context.claimToken,
task: updatedTask,
order: context.order,
orderItem: context.orderItem,
session,
})
}
async function finalizeClaimTaskRedeem(context) {
if (!context.task.browser_session_id) {
throw createHttpError('当前任务还没有创建浏览器会话', {
statusCode: 409,
errorCode: 'claim_session_not_created',
})
}
if (context.task.task_status !== 'role_confirmed' && context.task.task_status !== 'redeeming') {
throw createHttpError('当前任务还未确认角色,不能开始兑换', {
statusCode: 409,
errorCode: 'claim_role_not_confirmed',
})
}
const inventoryItem = context.task.primary_inventory_item_id
? await getInventoryItemById(context.task.primary_inventory_item_id)
: null
if (
!inventoryItem ||
String(inventoryItem.status || '').trim() !== 'reserved' ||
!String(inventoryItem.display_value || '').trim()
) {
throw createHttpError('当前任务没有可用的预占库存凭据', {
statusCode: 409,
errorCode: 'claim_inventory_not_reserved',
})
}
await updateTask(context.task.id, {
task_status: 'redeeming',
delivery_status: 'processing',
updated_at: nowIso(),
last_error: '',
})
try {
const redeemResult = await redeemClaimTaskWithInventoryFallback(context, inventoryItem)
const { session, inventoryItem: deliveredInventoryItem, classification, attempts } = redeemResult
const finalRedeem = session.redeem?.final?.redeem || null
const finishedAt = nowIso()
const updatedTask = await updateTask(context.task.id, {
task_status: 'redeemed',
inventory_status: 'consumed',
delivery_status: 'delivered',
result_code: resolveTencentRedeemResultCode(finalRedeem, classification),
result_message: classification.message,
screenshot_path: session.artifacts?.hasScreenshot ? await getTencentBrowserSessionScreenshotPath(session.sessionId) : '',
artifacts_json: JSON.stringify(session.artifacts || {}),
context_json: JSON.stringify(mergeTaskContext(context.task, {
redeemResolution: {
status: 'success',
attempts,
replacementCount: Math.max(0, attempts.length - 1),
finishedAt,
},
})),
redeemed_at: finishedAt,
updated_at: finishedAt,
last_error: '',
})
await markInventoryItemDelivered(deliveredInventoryItem.id, finishedAt)
await createTaskEvent(context.task.id, 'claim_redeem_completed', {
inventoryItemId: deliveredInventoryItem.id,
codeMasked: maskCode(deliveredInventoryItem.display_value),
replacementCount: Math.max(0, attempts.length - 1),
resultCode: resolveTencentRedeemResultCode(finalRedeem, classification),
resultMessage: classification.message,
}, finishedAt)
const autoDeliveryResult = await ensureAgisoXianyuAutoDeliveryForDeliveredTask({
order: context.order,
task: updatedTask,
trigger: 'claim_redeemed',
})
return buildClaimDetailPayload({
claimToken: context.claimToken,
task: autoDeliveryResult.task || updatedTask,
order: context.order,
orderItem: context.orderItem,
session,
})
} catch (error) {
const now = nowIso()
const nextRetryCount = Number(context.task.attempt_count || 0) + 1
const failureState = (
error && typeof error === 'object'
? /** @type {{ redeemTaskState?: { taskStatus: string, inventoryStatus: string, deliveryStatus: string, lastError: string, attempts: unknown[], classification: { retCode?: unknown } | null } }} */ (error).redeemTaskState
: null
) || {
taskStatus: 'retry_pending',
inventoryStatus: inventoryItem ? 'reserved' : 'pending',
deliveryStatus: 'pending',
lastError: error instanceof Error ? error.message : String(error || ''),
attempts: [],
classification: null,
}
const failedTask = await updateTask(context.task.id, {
task_status: failureState.taskStatus,
inventory_status: failureState.inventoryStatus,
delivery_status: failureState.deliveryStatus,
result_code: failureState.classification?.retCode != null ? String(failureState.classification.retCode) : '',
result_message: String(failureState.lastError || ''),
attempt_count: nextRetryCount,
last_error: String(failureState.lastError || ''),
context_json: JSON.stringify(mergeTaskContext(context.task, {
redeemResolution: {
status: 'failed',
taskStatus: failureState.taskStatus,
attempts: failureState.attempts || [],
replacementCount: Math.max(0, Number(failureState.attempts?.length || 1) - 1),
finishedAt: now,
},
})),
updated_at: now,
})
if (['retry_pending', 'waiting_inventory', 'manual_review'].includes(String(failureState.taskStatus || '').trim())) {
await notifyClaimRedeemNeedsAttention({
task: failedTask,
order: context.order,
status: failureState.taskStatus,
errorMessage: failureState.lastError,
})
}
throw Object.assign(error instanceof Error ? error : new Error(String(error || '兑换失败')), {
task: failedTask,
})
}
}
@@ -1,150 +0,0 @@
import { findClaimTokenByToken, getClaimTokenById, updateClaimToken } from '../../../repositories/claim-token-repo.js'
import { getOrderById } from '../../../repositories/order-repo.js'
import { getOrderItemById } from '../../../repositories/order-item-repo.js'
import { findTaskByClaimTokenId, getTaskById, updateTask } from '../../../repositories/task-repo.js'
import { releaseReservedInventoryItem } from '../../../repositories/inventory-repo.js'
import { createHttpError } from '../../../utils/http.js'
import { nowIso } from '../../../utils/time.js'
import { CLAIM_TERMINAL_STATUSES } from './shared.js'
export async function getClaimContext(token) {
const normalized = String(token || '').trim()
if (!normalized) {
throw createHttpError('缺少领取 token', {
statusCode: 400,
errorCode: 'missing_claim_token',
})
}
const claimToken = await findClaimTokenByToken(normalized)
if (!claimToken) {
throw createHttpError('领取链接无效或不存在', {
statusCode: 404,
errorCode: 'claim_token_not_found',
})
}
const task = await findTaskByClaimTokenId(claimToken.id)
if (!task) {
throw createHttpError('领取任务不存在', {
statusCode: 404,
errorCode: 'claim_task_not_found',
})
}
if (claimToken.status !== 'active') {
throw createHttpError('领取链接当前不可用', {
statusCode: 410,
errorCode: 'claim_token_inactive',
})
}
if (claimToken.expired_at && new Date(claimToken.expired_at).getTime() <= Date.now()) {
const expiredContext = await expireClaimContext(claimToken, task)
throw createHttpError('领取链接已过期', {
statusCode: 410,
errorCode: 'claim_token_expired',
context: expiredContext,
})
}
const [order, orderItem] = await Promise.all([
getOrderById(task.order_id),
getOrderItemById(task.order_item_id),
])
if (!order || !orderItem) {
throw createHttpError('领取任务关联订单不完整', {
statusCode: 500,
errorCode: 'claim_order_incomplete',
})
}
return {
claimToken,
task,
order,
orderItem,
}
}
export async function getClaimContextByTaskId(taskId) {
const task = await getTaskById(Number(taskId))
if (!task) {
throw createHttpError('领取任务不存在', {
statusCode: 404,
errorCode: 'claim_task_not_found',
})
}
const claimTokenId = Number(task.primary_claim_token_id || 0)
if (!claimTokenId) {
throw createHttpError('当前任务还没有领取链接', {
statusCode: 409,
errorCode: 'claim_token_missing',
})
}
const claimToken = await getClaimTokenById(claimTokenId)
if (!claimToken) {
throw createHttpError('领取链接不存在', {
statusCode: 404,
errorCode: 'claim_token_not_found',
})
}
const [order, orderItem] = await Promise.all([
getOrderById(task.order_id),
getOrderItemById(task.order_item_id),
])
if (!order || !orderItem) {
throw createHttpError('领取任务关联订单不完整', {
statusCode: 500,
errorCode: 'claim_order_incomplete',
})
}
return {
claimToken,
task,
order,
orderItem,
}
}
async function expireClaimContext(claimToken, task) {
const now = nowIso()
const nextClaimToken = await updateClaimToken(claimToken.id, {
status: 'expired',
updated_at: now,
})
let nextTask = task
if (!CLAIM_TERMINAL_STATUSES.has(String(task.task_status || '')) && task.task_status !== 'redeemed') {
if (task.primary_inventory_item_id) {
await releaseReservedInventoryItem(task.primary_inventory_item_id, now)
}
nextTask = await updateTask(task.id, {
task_status: 'expired',
inventory_status: 'pending',
user_action_status: 'expired',
last_error: '领取链接已过期,预占库存项已释放',
updated_at: now,
})
}
return {
claimToken: nextClaimToken,
task: nextTask,
}
}
@@ -1,175 +0,0 @@
import {
closeTencentBrowserSession,
createTencentBrowserSession,
reloadTencentBrowserSession,
} from '../../session/session.js'
import { updateTask } from '../../../repositories/task-repo.js'
import { nowIso } from '../../../utils/time.js'
import { getClaimContext, getClaimContextByTaskId } from './context.js'
import { clearTaskSession, loadTaskSession, syncTaskWithSession } from './runtime.js'
import {
assertTaskCanProceed,
buildClaimDetailPayload,
isRecoverableSessionError,
normalizeClaimLoginType,
} from './shared.js'
type JsonObject = Record<string, any>
export async function createClaimSession(token: unknown, payload: JsonObject = {}) {
return createClaimSessionWithContextLoader(
() => getClaimContext(token),
payload,
)
}
export async function createClaimSessionForAdminTask(taskId: unknown, payload: JsonObject = {}) {
return createClaimSessionWithContextLoader(
() => getClaimContextByTaskId(taskId),
payload,
)
}
export async function reloadClaimSession(token: unknown) {
return reloadClaimSessionWithContextLoader(() => getClaimContext(token))
}
export async function reloadClaimSessionForAdminTask(taskId: unknown) {
return reloadClaimSessionWithContextLoader(() => getClaimContextByTaskId(taskId))
}
export async function closeClaimSession(token: unknown) {
return closeClaimSessionWithContextLoader(() => getClaimContext(token))
}
export async function closeClaimSessionForAdminTask(taskId: unknown) {
return closeClaimSessionWithContextLoader(() => getClaimContextByTaskId(taskId))
}
async function createClaimSessionWithContextLoader(loadContext: () => Promise<JsonObject>, payload: JsonObject = {}) {
const context = await loadContext()
assertTaskCanProceed(context.task)
const requestedLoginType = normalizeClaimLoginType(payload.loginType)
const forceRecreate = Boolean(payload.forceRecreate)
let task = context.task
if (task.browser_session_id) {
const existing = await loadTaskSession(task, { includeQrImage: true })
task = existing.task
if (existing.session) {
const existingLoginType = normalizeClaimLoginType(existing.session.loginType || task.login_type)
if (!forceRecreate && existingLoginType === requestedLoginType) {
const syncedTask = await syncTaskWithSession(task, existing.session)
return buildClaimDetailPayload({
claimToken: context.claimToken,
task: syncedTask,
order: context.order,
orderItem: context.orderItem,
session: existing.session,
})
}
try {
await closeTencentBrowserSession(existing.session.sessionId)
} catch (error) {
if (!isRecoverableSessionError(error)) {
throw error
}
}
task = await clearTaskSession(task, {
lastError: '',
})
}
}
const session = await createTencentBrowserSession({
loginType: requestedLoginType,
})
const updatedTask = await updateTask(task.id, {
task_status: 'claimed',
browser_session_id: session.sessionId,
login_type: session.loginType,
user_action_status: 'claimed',
claimed_at: task.claimed_at || nowIso(),
updated_at: nowIso(),
last_error: '',
})
return buildClaimDetailPayload({
claimToken: context.claimToken,
task: updatedTask,
order: context.order,
orderItem: context.orderItem,
session,
})
}
async function reloadClaimSessionWithContextLoader(loadContext: () => Promise<JsonObject>) {
const context = await loadContext()
assertTaskCanProceed(context.task)
const active = await loadTaskSession(context.task, {
includeQrImage: true,
})
if (!active.session) {
return buildClaimDetailPayload({
claimToken: context.claimToken,
task: active.task,
order: context.order,
orderItem: context.orderItem,
session: null,
})
}
const session = await reloadTencentBrowserSession(active.session.sessionId)
const syncedTask = await syncTaskWithSession(active.task, session)
return buildClaimDetailPayload({
claimToken: context.claimToken,
task: syncedTask,
order: context.order,
orderItem: context.orderItem,
session,
})
}
async function closeClaimSessionWithContextLoader(loadContext: () => Promise<JsonObject>) {
const context = await loadContext()
assertTaskCanProceed(context.task)
let task = context.task
if (!task.browser_session_id) {
return buildClaimDetailPayload({
claimToken: context.claimToken,
task,
order: context.order,
orderItem: context.orderItem,
session: null,
})
}
try {
await closeTencentBrowserSession(task.browser_session_id)
} catch (error) {
if (!isRecoverableSessionError(error)) {
throw error
}
}
task = await clearTaskSession(task, {
lastError: '',
})
return buildClaimDetailPayload({
claimToken: context.claimToken,
task,
order: context.order,
orderItem: context.orderItem,
session: null,
})
}
@@ -1,301 +0,0 @@
import { reloadTencentBrowserSession, redeemTencentBrowserSession } from '../../session/session.js'
import { classifyTencentRedeemResult } from '../../session/session-redeem.js'
import { createTaskEvent } from '../../../repositories/task-event-repo.js'
import {
invalidateReservedInventoryItem,
markInventoryItemConsumed,
} from '../../../repositories/inventory-repo.js'
import { reserveInventoryForTask } from '../../order/inventory-service.js'
import { createHttpError } from '../../../utils/http.js'
import { nowIso } from '../../../utils/time.js'
import { maskCode, REDEEM_REPLACEMENT_LIMIT } from './shared.js'
import type { InventoryItemRow, OrderItemRow, TaskRow } from '../../../types/repository-rows.js'
type ClaimRedeemContext = {
task: Pick<TaskRow, 'id' | 'browser_session_id'>
orderItem: Pick<OrderItemRow, 'sku_code'>
}
type TencentRedeemPayload = {
[key: string]: unknown
iRet?: unknown
ret?: unknown
}
type TencentRedeemSession = {
[key: string]: unknown
sessionId?: string
artifacts?: {
[key: string]: unknown
hasScreenshot?: boolean
} | null
redeem?: {
final?: {
redeem?: TencentRedeemPayload | null
} | null
} | null
}
type TencentRedeemClassification = {
success: boolean
outcome: string
message: string
retCode?: unknown
[key: string]: unknown
}
type RedeemAttemptSummary = {
attempt: number
inventoryItemId: number
codeMasked: string
credentialType: string
outcome: string
resultCode: string
resultMessage: string
}
type RedeemTaskState = {
taskStatus: string
inventoryStatus: string
deliveryStatus: string
lastError: string
attempts: RedeemAttemptSummary[]
classification: TencentRedeemClassification | null
}
type RedeemTaskStateError = Error & {
statusCode?: number
errorCode?: string
redeemTaskState: RedeemTaskState
}
type RedeemTaskStateErrorPayload = {
statusCode?: number
errorCode?: string
taskStatus?: string
inventoryStatus?: string
deliveryStatus?: string
attempts?: RedeemAttemptSummary[]
classification?: TencentRedeemClassification | null
}
type RedeemClaimTaskDeps = {
reloadTencentBrowserSession?: (sessionId: string) => Promise<unknown>
redeemTencentBrowserSession?: (
sessionId: string,
payload: { code: string },
) => Promise<TencentRedeemSession>
classifyTencentRedeemResult?: (finalRedeem: TencentRedeemPayload | null) => TencentRedeemClassification
markInventoryItemConsumed?: (
inventoryItemId: number | string,
reason: string,
consumedAt: string,
) => Promise<InventoryItemRow | null>
createTaskEvent?: (
taskId: number | string,
eventType: string,
payload?: unknown,
createdAt?: string,
) => Promise<unknown>
invalidateReservedInventoryItem?: (
inventoryItemId: number | string,
reason: string,
updatedAt: string,
) => Promise<InventoryItemRow | null>
reserveInventoryForTask?: (payload: {
skuCode: string
taskId: number | string
credentialType?: string
roleKey?: string
inventoryGroupCodes?: string[] | null
}) => Promise<InventoryItemRow | null>
nowIso?: () => string
redeemReplacementLimit?: number
}
type RedeemClaimTaskResult = {
session: TencentRedeemSession
inventoryItem: InventoryItemRow
classification: TencentRedeemClassification
attempts: RedeemAttemptSummary[]
}
export async function redeemClaimTaskWithInventoryFallback(
context: ClaimRedeemContext,
initialInventoryItem: InventoryItemRow,
): Promise<RedeemClaimTaskResult> {
return redeemClaimTaskWithInventoryFallbackWithDeps(context, initialInventoryItem)
}
export async function redeemClaimTaskWithInventoryFallbackWithDeps(
context: ClaimRedeemContext,
initialInventoryItem: InventoryItemRow,
{
reloadTencentBrowserSession: reloadRedeemSession = reloadTencentBrowserSession,
redeemTencentBrowserSession: redeemSession = redeemTencentBrowserSession,
classifyTencentRedeemResult: classifyRedeemResult = classifyTencentRedeemResult,
markInventoryItemConsumed: markConsumedInventoryItem = markInventoryItemConsumed,
createTaskEvent: createRedeemTaskEvent = createTaskEvent,
invalidateReservedInventoryItem: invalidateReservedItem = invalidateReservedInventoryItem,
reserveInventoryForTask: reserveReplacementInventory = reserveInventoryForTask,
nowIso: getNowIso = nowIso,
redeemReplacementLimit = REDEEM_REPLACEMENT_LIMIT,
}: RedeemClaimTaskDeps = {},
): Promise<RedeemClaimTaskResult> {
const attempts: RedeemAttemptSummary[] = []
let currentInventoryItem = initialInventoryItem
let shouldReloadBeforeNextAttempt = false
for (let attemptIndex = 1; attemptIndex <= redeemReplacementLimit; attemptIndex += 1) {
if (shouldReloadBeforeNextAttempt) {
await reloadRedeemSession(context.task.browser_session_id)
shouldReloadBeforeNextAttempt = false
}
const session = await redeemSession(context.task.browser_session_id, {
code: currentInventoryItem.display_value,
})
const finalRedeem = session.redeem?.final?.redeem || null
const classification = classifyRedeemResult(finalRedeem)
const attemptSummary = {
attempt: attemptIndex,
inventoryItemId: currentInventoryItem.id,
codeMasked: maskCode(currentInventoryItem.display_value),
credentialType: String(currentInventoryItem.credential_type || ''),
outcome: classification.outcome,
resultCode: resolveTencentRedeemResultCode(finalRedeem, classification),
resultMessage: classification.message,
}
attempts.push(attemptSummary)
if (classification.success) {
return {
session,
inventoryItem: currentInventoryItem,
classification,
attempts,
}
}
if (classification.outcome === 'code_used') {
const updatedAt = getNowIso()
await markConsumedInventoryItem(currentInventoryItem.id, classification.message, updatedAt)
await createRedeemTaskEvent(context.task.id, 'claim_redeem_code_used', {
inventoryItemId: currentInventoryItem.id,
codeMasked: attemptSummary.codeMasked,
resultCode: attemptSummary.resultCode,
resultMessage: classification.message,
}, updatedAt)
} else if (classification.outcome === 'code_invalid') {
const updatedAt = getNowIso()
await invalidateReservedItem(currentInventoryItem.id, classification.message, updatedAt)
await createRedeemTaskEvent(context.task.id, 'claim_redeem_code_invalid', {
inventoryItemId: currentInventoryItem.id,
codeMasked: attemptSummary.codeMasked,
resultCode: attemptSummary.resultCode,
resultMessage: classification.message,
}, updatedAt)
} else {
throw createRedeemTaskStateError(classification.message, {
errorCode: 'claim_redeem_failed',
taskStatus: 'retry_pending',
inventoryStatus: 'reserved',
deliveryStatus: 'pending',
attempts,
classification,
})
}
if (attemptIndex >= redeemReplacementLimit) {
throw createRedeemTaskStateError('连续更换兑换码后仍未成功,请联系人工处理', {
errorCode: 'claim_redeem_replacement_limit_reached',
taskStatus: 'retry_pending',
inventoryStatus: 'pending',
deliveryStatus: 'pending',
attempts,
classification,
})
}
const replacement = await reserveReplacementInventory({
skuCode: context.orderItem.sku_code,
taskId: context.task.id,
credentialType: currentInventoryItem.credential_type || 'tencent_code',
roleKey: 'primary_code',
inventoryGroupCodes: currentInventoryItem.inventory_group_code
? [String(currentInventoryItem.inventory_group_code).trim()]
: null,
})
if (!replacement || !String(replacement.display_value || '').trim()) {
const exhaustedMessage = classification.outcome === 'code_used'
? '兑换码已使用,且没有更多同类型可用 CDK 可继续重试'
: '兑换码错误,请确认库存数据;当前没有更多同类型可用 CDK 可继续重试'
throw createRedeemTaskStateError(exhaustedMessage, {
errorCode: 'claim_inventory_replacement_exhausted',
taskStatus: 'waiting_inventory',
inventoryStatus: 'pending',
deliveryStatus: 'pending',
attempts,
classification,
})
}
const replacedAt = getNowIso()
await createRedeemTaskEvent(context.task.id, 'claim_redeem_inventory_replaced', {
previousInventoryItemId: currentInventoryItem.id,
previousCodeMasked: attemptSummary.codeMasked,
previousOutcome: classification.outcome,
nextInventoryItemId: replacement.id,
nextCodeMasked: maskCode(replacement.display_value),
credentialType: String(replacement.credential_type || currentInventoryItem.credential_type || ''),
inventoryGroupCode: String(replacement.inventory_group_code || currentInventoryItem.inventory_group_code || '').trim(),
}, replacedAt)
currentInventoryItem = replacement
shouldReloadBeforeNextAttempt = true
}
throw createRedeemTaskStateError('兑换失败,请稍后重试', {
errorCode: 'claim_redeem_failed',
taskStatus: 'retry_pending',
inventoryStatus: 'pending',
deliveryStatus: 'pending',
attempts,
classification: null,
})
}
function createRedeemTaskStateError(
message: string,
payload: RedeemTaskStateErrorPayload = {},
): RedeemTaskStateError {
const error = createHttpError(message, {
statusCode: payload.statusCode || 409,
errorCode: payload.errorCode || 'claim_redeem_failed',
}) as unknown as RedeemTaskStateError
error.redeemTaskState = {
taskStatus: payload.taskStatus || 'retry_pending',
inventoryStatus: payload.inventoryStatus || 'pending',
deliveryStatus: payload.deliveryStatus || 'pending',
lastError: String(message || ''),
attempts: payload.attempts || [],
classification: payload.classification || null,
}
return error
}
export function resolveTencentRedeemResultCode(
finalRedeem: TencentRedeemPayload | null,
classification: Pick<TencentRedeemClassification, 'retCode'> | null,
): string {
if (classification?.retCode != null) {
return String(classification.retCode)
}
const value = finalRedeem?.iRet ?? finalRedeem?.ret
return value == null ? '' : String(value)
}
@@ -1,107 +0,0 @@
import {
getTencentBrowserSession,
getTencentBrowserSessionSummary,
} from '../../session/session.js'
import { updateTask } from '../../../repositories/task-repo.js'
import { nowIso } from '../../../utils/time.js'
import { isRecoverableSessionError, parseTaskState } from './shared.js'
type JsonObject = Record<string, any>
export async function loadTaskSession(task: JsonObject, { includeQrImage = false }: { includeQrImage?: boolean } = {}) {
if (!task.browser_session_id) {
return {
task,
session: null,
}
}
try {
const session = includeQrImage
? await getTencentBrowserSession(task.browser_session_id)
: await getTencentBrowserSessionSummary(task.browser_session_id)
return {
task,
session,
}
} catch (error) {
if (!isRecoverableSessionError(error)) {
throw error
}
const nextTask = await clearTaskSession(task, {
lastError: '浏览器会话已失效,请重新初始化登录',
})
return {
task: nextTask,
session: null,
}
}
}
export async function syncTaskWithSession(task: JsonObject, session: JsonObject) {
const activityInfo = session.activityInfo || null
const patch: JsonObject = {
login_type: String(session.loginType || task.login_type || ''),
updated_at: nowIso(),
}
if (activityInfo?.nickname) {
patch.nickname = String(activityInfo.nickname)
}
if (activityInfo?.role?.ready) {
patch.role_id = String(activityInfo.role.roleId || '')
patch.role_name = String(activityInfo.role.roleName || '')
patch.area = String(activityInfo.role.area || '')
patch.partition_name = String(activityInfo.role.partition || '')
}
if (task.task_status === 'link_generated') {
patch.task_status = 'claimed'
patch.user_action_status = 'claimed'
patch.claimed_at = task.claimed_at || nowIso()
}
if (session.review?.capturedAt) {
patch.state_json = JSON.stringify({
...parseTaskState(task),
reviewScreenshotReady: true,
reviewCapturedAt: String(session.review.capturedAt || ''),
reviewRoleId: String(session.review.roleId || ''),
reviewRoleName: String(session.review.roleName || ''),
})
}
if (session.status === 'redeemed' && session.artifacts?.hasScreenshot) {
patch.screenshot_path = task.screenshot_path || ''
}
return updateTask(task.id, patch)
}
export async function clearTaskSession(task: JsonObject, { lastError = '' }: { lastError?: string } = {}) {
const shouldResetClaimProgress = ['link_generated', 'claimed', 'role_confirmed'].includes(String(task.task_status || ''))
const nextTaskStatus = shouldResetClaimProgress ? 'link_generated' : task.task_status
const nextUserActionStatus = shouldResetClaimProgress ? 'pending_claim' : task.user_action_status
const patch = {
task_status: nextTaskStatus,
user_action_status: nextUserActionStatus,
browser_session_id: '',
login_type: '',
nickname: '',
role_id: '',
role_name: '',
area: '',
partition_name: '',
artifacts_json: '{}',
state_json: '{}',
role_confirmed_at: nextTaskStatus === 'link_generated' ? null : task.role_confirmed_at,
last_error: String(lastError || ''),
updated_at: nowIso(),
}
return updateTask(task.id, patch)
}