From b35938c14820325aececa6b9a797cc58250cc99a Mon Sep 17 00:00:00 2001 From: yml Date: Wed, 20 May 2026 00:38:18 +0800 Subject: [PATCH] =?UTF-8?q?feat(cloudtentacles):=20=E5=90=8E=E7=AB=AF?= =?UTF-8?q?=E6=A0=B8=E5=BF=83=E6=95=B0=E6=8D=AE=E7=BB=93=E6=9E=84=E5=A4=9A?= =?UTF-8?q?=E8=B4=A6=E5=8F=B7=E6=94=AF=E6=8C=81=20-=20source-config/sessio?= =?UTF-8?q?n-state=20=E6=94=B9=E9=80=A0=20+=20resolvePersistedCloudtentacl?= =?UTF-8?q?esContext=20=E5=A2=9E=E5=8A=A0=20sourceKey=20=E5=8F=82=E6=95=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../kuaishou-cloud-task-service.js | 1335 ++++++++++------- .../cloudtentacles/session-state-service.js | 114 +- .../cloudtentacles/source-config-service.js | 148 +- 3 files changed, 1022 insertions(+), 575 deletions(-) diff --git a/apps/backend/src/services/fulfillment/kuaishou-cloud-task-service.js b/apps/backend/src/services/fulfillment/kuaishou-cloud-task-service.js index 6c17b7d5..c2e4943c 100644 --- a/apps/backend/src/services/fulfillment/kuaishou-cloud-task-service.js +++ b/apps/backend/src/services/fulfillment/kuaishou-cloud-task-service.js @@ -1,19 +1,25 @@ // @ts-check -import { createTaskEvent } from '../../repositories/task-event-repo.js' -import { getOrderById } from '../../repositories/order-repo.js' -import { updateTask } from '../../repositories/task-repo.js' -import { buildClaimUrl, createTaskClaimToken } from '../claim/claim-service.js' -import { normalizeProductName } from '../order/product-match-service.js' -import { getCloudtentaclesSourceConfig } from '../platforms/cloudtentacles/source-config-service.js' -import { getCloudtentaclesSessionState } from '../platforms/cloudtentacles/session-state-service.js' +import { createTaskEvent } from "../../repositories/task-event-repo.js"; +import { getOrderById } from "../../repositories/order-repo.js"; +import { updateTask } from "../../repositories/task-repo.js"; +import { buildClaimUrl, createTaskClaimToken } from "../claim/claim-service.js"; +import { normalizeProductName } from "../order/product-match-service.js"; +import { + getCloudtentaclesSourceConfig, + getCloudtentaclesSourceByKey, +} from "../platforms/cloudtentacles/source-config-service.js"; +import { + getCloudtentaclesSessionState, + getCloudtentaclesSessionStateByKey, +} 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' +} from "../platforms/cloudtentacles/catalog-service.js"; +import { getCloudtentaclesKnapsack } from "../platforms/cloudtentacles/knapsack-service.js"; import { appointCloudtentaclesVirtualNumber, backCloudtentaclesVirtualNumber, @@ -23,82 +29,105 @@ import { getCloudtentaclesBindUrl, probeCloudtentaclesBindUrl, verifyCloudtentaclesLoginCode, -} from '../platforms/cloudtentacles/virtual-number-service.js' -import { resolveCloudtentaclesConfig } from '../platforms/cloudtentacles/shared.js' -import { consumeKuaishouEticket } from '../platforms/kuaishou-eticket/consume-service.js' +} from "../platforms/cloudtentacles/virtual-number-service.js"; +import { resolveCloudtentaclesConfig } from "../platforms/cloudtentacles/shared.js"; +import { consumeKuaishouEticket } from "../platforms/kuaishou-eticket/consume-service.js"; import { getKuaishouEticketSourceConfig, resolveKuaishouEticketShopConfig, -} from '../platforms/kuaishou-eticket/source-config-service.js' +} from "../platforms/kuaishou-eticket/source-config-service.js"; import { notifyKuaishouCloudAssetNotEnough, notifyKuaishouCloudBindUrlRefreshFailed, notifyKuaishouCloudConsumeFailed, -} from '../notification/domain-notifications.js' -import { createHttpError } from '../../utils/http.js' -import { nowIso } from '../../utils/time.js' +} from "../notification/domain-notifications.js"; +import { createHttpError } from "../../utils/http.js"; +import { nowIso } from "../../utils/time.js"; -const KUAISHOU_CLOUD_FIXED_VN_KEY = '1' +const KUAISHOU_CLOUD_FIXED_VN_KEY = "1"; export function isKuaishouCloudTask(task) { - return String(task?.executor_key || '').trim() === 'kuaishou_ct_assisted' + return String(task?.executor_key || "").trim() === "kuaishou_ct_assisted"; } export function normalizeKuaishouCloudFlow(value) { - const source = value && typeof value === 'object' ? value : {} - const binding = source.binding && typeof source.binding === 'object' ? source.binding : {} - const role = source.role && typeof source.role === 'object' ? source.role : {} - const purchase = source.purchase && typeof source.purchase === 'object' ? source.purchase : {} - const dispatch = source.dispatch && typeof source.dispatch === 'object' ? source.dispatch : {} - const returnNumber = source.returnNumber && typeof source.returnNumber === 'object' ? source.returnNumber : {} - const consume = source.consume && typeof source.consume === 'object' ? source.consume : {} - const ticket = source.ticket && typeof source.ticket === 'object' ? source.ticket : {} + const source = value && typeof value === "object" ? value : {}; + const binding = + source.binding && typeof source.binding === "object" ? source.binding : {}; + const role = + source.role && typeof source.role === "object" ? source.role : {}; + const purchase = + source.purchase && typeof source.purchase === "object" + ? source.purchase + : {}; + const dispatch = + source.dispatch && typeof source.dispatch === "object" + ? source.dispatch + : {}; + const returnNumber = + source.returnNumber && typeof source.returnNumber === "object" + ? source.returnNumber + : {}; + const consume = + source.consume && typeof source.consume === "object" ? source.consume : {}; + const ticket = + source.ticket && typeof source.ticket === "object" ? source.ticket : {}; - const roleName = String(role.name || binding.roleName || '').trim() - const roleId = String(role.rid || binding.roleId || '').trim() - const bindPreparedAt = binding.bindPreparedAt || null - const bindExpiresAt = binding.bindExpiresAt || resolveKuaishouCloudBindUrlExpiresAt(bindPreparedAt) + const roleName = String(role.name || binding.roleName || "").trim(); + const roleId = String(role.rid || binding.roleId || "").trim(); + const bindPreparedAt = binding.bindPreparedAt || null; + const bindExpiresAt = + binding.bindExpiresAt || + resolveKuaishouCloudBindUrlExpiresAt(bindPreparedAt); return { ...source, - configId: String(source.configId || '').trim(), - internalSkuCode: String(source.internalSkuCode || '').trim(), - internalSkuName: String(source.internalSkuName || '').trim(), + configId: String(source.configId || "").trim(), + internalSkuCode: String(source.internalSkuCode || "").trim(), + internalSkuName: String(source.internalSkuName || "").trim(), ticket: { - code: String(ticket.code || '').trim(), - status: String(ticket.status || 'pending').trim() || 'pending', + code: String(ticket.code || "").trim(), + status: String(ticket.status || "pending").trim() || "pending", capturedAt: ticket.capturedAt || null, capturedBy: ticket.capturedBy || null, verifiedAt: ticket.verifiedAt || null, - oid: String(ticket.oid || '').trim(), - formToken: String(ticket.formToken || '').trim(), + oid: String(ticket.oid || "").trim(), + formToken: String(ticket.formToken || "").trim(), leftCount: Number(ticket.leftCount || 0) || 0, - goodsTitle: String(ticket.goodsTitle || '').trim(), + goodsTitle: String(ticket.goodsTitle || "").trim(), }, binding: { - prepareStatus: String(binding.prepareStatus || 'pending').trim() || 'pending', - cloudSourceKey: String(binding.cloudSourceKey || 'default').trim() || 'default', + 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 || KUAISHOU_CLOUD_FIXED_VN_KEY).trim() || KUAISHOU_CLOUD_FIXED_VN_KEY, + skuName: String(binding.skuName || "").trim(), + vnKey: + String(binding.vnKey || KUAISHOU_CLOUD_FIXED_VN_KEY).trim() || + KUAISHOU_CLOUD_FIXED_VN_KEY, vnId: Number(binding.vnId || 0) || 0, - vnPhone: String(binding.vnPhone || '').trim(), - bindUrl: String(binding.bindUrl || '').trim(), + vnPhone: String(binding.vnPhone || "").trim(), + bindUrl: String(binding.bindUrl || "").trim(), bindPreparedAt, bindExpiresAt, bindProbeAt: binding.bindProbeAt || null, - bindProbeStatus: String(binding.bindProbeStatus || '').trim(), - bindProbeMessage: String(binding.bindProbeMessage || '').trim(), + bindProbeStatus: String(binding.bindProbeStatus || "").trim(), + bindProbeMessage: String(binding.bindProbeMessage || "").trim(), roleName, roleId, }, role: { - status: String(role.status || (roleName || roleId ? 'ready' : 'pending')).trim() || 'pending', + status: + String( + role.status || (roleName || roleId ? "ready" : "pending") + ).trim() || "pending", name: roleName, rid: roleId, refreshedAt: role.refreshedAt || null, - errorMessage: String(role.errorMessage || '').trim(), - rawInfo: role.rawInfo && typeof role.rawInfo === 'object' ? role.rawInfo : null, + errorMessage: String(role.errorMessage || "").trim(), + rawInfo: + role.rawInfo && typeof role.rawInfo === "object" ? role.rawInfo : null, }, purchase: { autoBuyEnabled: purchase.autoBuyEnabled !== false, @@ -110,116 +139,149 @@ export function normalizeKuaishouCloudFlow(value) { purchaseAt: purchase.purchaseAt || null, }, dispatch: { - status: String(dispatch.status || 'pending').trim() || 'pending', + 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(), + note: String(dispatch.note || "").trim(), }, returnNumber: { - status: String(returnNumber.status || 'pending').trim() || 'pending', + status: String(returnNumber.status || "pending").trim() || "pending", returnedAt: returnNumber.returnedAt || null, returnedBy: returnNumber.returnedBy || null, autoReturnEnabled: returnNumber.autoReturnEnabled === true, }, consume: { - status: String(consume.status || 'pending').trim() || 'pending', - shopId: String(consume.shopId || '').trim(), - shopName: String(consume.shopName || '').trim(), + status: String(consume.status || "pending").trim() || "pending", + shopId: String(consume.shopId || "").trim(), + shopName: String(consume.shopName || "").trim(), autoConsumeEnabled: consume.autoConsumeEnabled === true, consumedAt: consume.consumedAt || null, - errorMessage: String(consume.errorMessage || '').trim(), + errorMessage: String(consume.errorMessage || "").trim(), }, - notes: String(source.notes || '').trim(), - } + notes: String(source.notes || "").trim(), + }; } export function normalizeKuaishouCloudRoleInfo(value) { - const rawInfo = value && typeof value === 'object' ? value : null - const nestedBindInfo = rawInfo?.sBindInfo && typeof rawInfo.sBindInfo === 'object' ? rawInfo.sBindInfo : null - const source = nestedBindInfo || rawInfo + const rawInfo = value && typeof value === "object" ? value : null; + const nestedBindInfo = + rawInfo?.sBindInfo && typeof rawInfo.sBindInfo === "object" + ? rawInfo.sBindInfo + : null; + const source = nestedBindInfo || rawInfo; return { - name: String(source?.name || source?.roleName || source?.nickname || source?.sRoleName || '').trim(), - rid: String(source?.rid || source?.roleId || source?.uid || source?.sRoleId || source?.sUserId || '').trim(), + name: String( + source?.name || + source?.roleName || + source?.nickname || + source?.sRoleName || + "" + ).trim(), + rid: String( + source?.rid || + source?.roleId || + source?.uid || + source?.sRoleId || + source?.sUserId || + "" + ).trim(), rawInfo, - } + }; } -export function resolvePersistedCloudtentaclesContext() { - const source = getCloudtentaclesSourceConfig() - const session = getCloudtentaclesSessionState() - const token = String(session.token || '').trim() +export function resolvePersistedCloudtentaclesContext(sourceKey = "default") { + const source = + getCloudtentaclesSourceByKey(sourceKey) || getCloudtentaclesSourceConfig(); + const session = + getCloudtentaclesSessionStateByKey(sourceKey) || + getCloudtentaclesSessionState(); + const token = String(session.token || "").trim(); if (!token) { - throw createHttpError('当前 cloudtentacles 没有可用 token,请先到平台配置完成登录校验', { - statusCode: 409, - errorCode: 'kuaishou_cloud_missing_cloud_token', - }) + throw createHttpError( + "当前 cloudtentacles 没有可用 token,请先到平台配置完成登录校验", + { + statusCode: 409, + errorCode: "kuaishou_cloud_missing_cloud_token", + } + ); } return { - baseUrl: String(session.baseUrl || source.baseUrl || '').trim() || 'https://123.207.217.176', + baseUrl: + String(session.baseUrl || source.baseUrl || "").trim() || + "https://123.207.217.176", token, - deviceId: String(session.deviceId || source.deviceId || '-').trim() || '-', + deviceId: String(session.deviceId || source.deviceId || "-").trim() || "-", deviceType: Number(session.deviceType ?? source.deviceType ?? 0), - } + }; } export async function ensureTaskClaimLink(task) { - const tokenStatus = String(task?.primary_claim_token_status || '').trim() - const token = String(task?.primary_claim_token || task?.claim_token || '').trim() - const expiredAt = getTaskClaimExpiresAt(task) + const tokenStatus = String(task?.primary_claim_token_status || "").trim(); + const token = String( + task?.primary_claim_token || task?.claim_token || "" + ).trim(); + const expiredAt = getTaskClaimExpiresAt(task); - if (tokenStatus === 'active' && token && !isClaimExpired(expiredAt)) { + if (tokenStatus === "active" && token && !isClaimExpired(expiredAt)) { return { token, expiredAt, claimUrl: buildClaimUrl(token), - } + }; } - const claimToken = await createTaskClaimToken(task.id) + const claimToken = await createTaskClaimToken(task.id); return { token: claimToken.token, expiredAt: claimToken.expired_at, claimUrl: claimToken.claimUrl, - } + }; } export async function prepareKuaishouCloudFulfillmentTask(task, options = {}) { if (!isKuaishouCloudTask(task)) { - throw createHttpError('当前任务不是快手 Cloud 履约任务', { + throw createHttpError("当前任务不是快手 Cloud 履约任务", { statusCode: 409, - errorCode: 'kuaishou_cloud_task_invalid', - }) + errorCode: "kuaishou_cloud_task_invalid", + }); } - const now = nowIso() - const actor = normalizeActor(options.actor) - const force = options.force === true - const taskContext = parseTaskContext(task) - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) - const claimLinkState = await ensureTaskClaimLink(task) + const now = nowIso(); + const actor = normalizeActor(options.actor); + const force = options.force === true; + const taskContext = parseTaskContext(task); + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment); + const claimLinkState = await ensureTaskClaimLink(task); - if (!force && flow.binding.prepareStatus === 'ready' && flow.binding.vnId > 0 && flow.binding.vnPhone && flow.binding.bindUrl) { + if ( + !force && + flow.binding.prepareStatus === "ready" && + flow.binding.vnId > 0 && + flow.binding.vnPhone && + flow.binding.bindUrl + ) { if (!isKuaishouCloudBindUrlFresh(flow)) { return refreshKuaishouCloudTaskBindUrl(task, { - source: options.source || 'system_refresh_expired_bind_url', + source: options.source || "system_refresh_expired_bind_url", actor, claimLinkState, - }) + }); } - let readyTask = task + let readyTask = task; - if (String(task.task_status || '').trim() !== 'waiting_binding') { + if (String(task.task_status || "").trim() !== "waiting_binding") { readyTask = await updateTask(task.id, { - task_status: 'waiting_binding', - claim_token: claimLinkState.token || task.claim_token || '', - claim_expires_at: claimLinkState.expiredAt || getTaskClaimExpiresAt(task), + task_status: "waiting_binding", + claim_token: claimLinkState.token || task.claim_token || "", + claim_expires_at: + claimLinkState.expiredAt || getTaskClaimExpiresAt(task), updated_at: now, - }) + }); } return { @@ -227,28 +289,30 @@ export async function prepareKuaishouCloudFulfillmentTask(task, options = {}) { claimUrl: claimLinkState.claimUrl, token: claimLinkState.token, flow, - } + }; } - const cloudContext = resolvePersistedCloudtentaclesContext() + const cloudContext = resolvePersistedCloudtentaclesContext( + flow.binding.cloudSourceKey || "default" + ); 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 资源,无法自动准备绑定', { + throw createHttpError("当前任务缺少可用 cloud 资源,无法自动准备绑定", { statusCode: 409, - errorCode: 'kuaishou_cloud_missing_binding_config', - }) + errorCode: "kuaishou_cloud_missing_binding_config", + }); } const flowWithResolvedBinding = { @@ -259,35 +323,39 @@ export async function prepareKuaishouCloudFulfillmentTask(task, options = {}) { skuName: resolvedBinding.skuName, vnKey: KUAISHOU_CLOUD_FIXED_VN_KEY, }, - } + }; - const knapsackItem = resolvedBinding.knapsackItem - let usedKnapsack = Number(knapsackItem?.count || 0) > 0 - let purchaseTriggered = false - let assetBefore = 0 - let assetAfter = 0 + const 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('背包中没有现成库存,且当前配置未开启自动购买', { + throw createHttpError("背包中没有现成库存,且当前配置未开启自动购买", { statusCode: 409, - errorCode: 'kuaishou_cloud_auto_buy_disabled', - }) + errorCode: "kuaishou_cloud_auto_buy_disabled", + }); } - const asset = await getCloudtentaclesAsset(cloudContext) - const targetSku = resolvedBinding.skuItem + const asset = await getCloudtentaclesAsset(cloudContext); + const targetSku = resolvedBinding.skuItem; if (!targetSku) { - throw createHttpError(`cloudtentacles 未找到 SKU ${flowWithResolvedBinding.binding.skuId}`, { - statusCode: 404, - errorCode: 'kuaishou_cloud_sku_not_found', - }) + throw createHttpError( + `cloudtentacles 未找到 SKU ${flowWithResolvedBinding.binding.skuId}`, + { + statusCode: 404, + errorCode: "kuaishou_cloud_sku_not_found", + } + ); } - assetBefore = Number(asset.asset || 0) || 0 - const targetPrice = Number(targetSku.price || 0) || 0 - const requiredAsset = targetPrice + flowWithResolvedBinding.purchase.minAssetReserve + assetBefore = Number(asset.asset || 0) || 0; + const targetPrice = Number(targetSku.price || 0) || 0; + const requiredAsset = + targetPrice + flowWithResolvedBinding.purchase.minAssetReserve; if (assetBefore < requiredAsset) { await notifyKuaishouCloudAssetNotEnough({ task, @@ -295,28 +363,31 @@ export async function prepareKuaishouCloudFulfillmentTask(task, options = {}) { assetBefore, requiredAsset, skuName: targetSku.name, - }) - throw createHttpError(`余额不足,当前 ${assetBefore},至少需要 ${requiredAsset}`, { - statusCode: 409, - errorCode: 'kuaishou_cloud_asset_not_enough', - }) + }); + throw createHttpError( + `余额不足,当前 ${assetBefore},至少需要 ${requiredAsset}`, + { + statusCode: 409, + errorCode: "kuaishou_cloud_asset_not_enough", + } + ); } await buyCloudtentaclesSku({ ...cloudContext, id: flowWithResolvedBinding.binding.skuId, count: 1, - }) - purchaseTriggered = true + }); + purchaseTriggered = true; - const assetResult = await getCloudtentaclesAsset(cloudContext) - assetAfter = Number(assetResult.asset || 0) || 0 + const assetResult = await getCloudtentaclesAsset(cloudContext); + assetAfter = Number(assetResult.asset || 0) || 0; } const preparedBinding = await prepareKuaishouCloudBindResourceWithFallback({ cloudContext, vnKeyCandidates, - }) + }); const nextContext = { ...taskContext, @@ -325,24 +396,24 @@ export async function prepareKuaishouCloudFulfillmentTask(task, options = {}) { binding: { ...flowWithResolvedBinding.binding, vnKey: preparedBinding.vnKey, - prepareStatus: 'ready', + prepareStatus: "ready", vnId: preparedBinding.vnId, vnPhone: preparedBinding.vnPhone, bindUrl: preparedBinding.bindUrl, bindPreparedAt: now, bindExpiresAt: resolveKuaishouCloudBindUrlExpiresAt(now), bindProbeAt: null, - bindProbeStatus: 'pending', - bindProbeMessage: '', - roleName: '', - roleId: '', + bindProbeStatus: "pending", + bindProbeMessage: "", + roleName: "", + roleId: "", }, role: { - status: 'pending', - name: '', - rid: '', + status: "pending", + name: "", + rid: "", refreshedAt: null, - errorMessage: '', + errorMessage: "", rawInfo: null, }, purchase: { @@ -351,85 +422,99 @@ export async function prepareKuaishouCloudFulfillmentTask(task, options = {}) { purchaseTriggered, assetBefore, assetAfter, - purchaseAt: purchaseTriggered ? now : flowWithResolvedBinding.purchase.purchaseAt, + purchaseAt: purchaseTriggered + ? now + : flowWithResolvedBinding.purchase.purchaseAt, }, }, - } + }; const updatedTask = await updateTask(task.id, { - task_status: 'waiting_binding', - inventory_status: 'not_required', - user_action_status: 'pending_claim', - claim_token: claimLinkState.token || task.claim_token || '', + task_status: "waiting_binding", + inventory_status: "not_required", + user_action_status: "pending_claim", + claim_token: claimLinkState.token || task.claim_token || "", claim_expires_at: claimLinkState.expiredAt || getTaskClaimExpiresAt(task), - role_id: '', - role_name: '', - last_error: '', + role_id: "", + role_name: "", + last_error: "", context_json: JSON.stringify(nextContext), updated_at: now, - }) + }); - await createTaskEvent(task.id, 'kuaishou_cloud_binding_prepared', { - source: String(options.source || 'system').trim() || 'system', - skuId: flowWithResolvedBinding.binding.skuId, - skuName: flowWithResolvedBinding.binding.skuName, - vnKey: preparedBinding.vnKey, - vnId: preparedBinding.vnId, - vnPhoneMasked: maskPhone(preparedBinding.vnPhone), - purchaseTriggered, - usedKnapsack, - resolvedByName: resolvedBinding.resolvedByName, - actor, - }, now) + await createTaskEvent( + task.id, + "kuaishou_cloud_binding_prepared", + { + source: String(options.source || "system").trim() || "system", + skuId: flowWithResolvedBinding.binding.skuId, + skuName: flowWithResolvedBinding.binding.skuName, + vnKey: preparedBinding.vnKey, + vnId: preparedBinding.vnId, + vnPhoneMasked: maskPhone(preparedBinding.vnPhone), + purchaseTriggered, + usedKnapsack, + resolvedByName: resolvedBinding.resolvedByName, + actor, + }, + now + ); return { task: updatedTask, claimUrl: claimLinkState.claimUrl, token: claimLinkState.token, flow: normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment), - } + }; } export async function refreshKuaishouCloudTaskBindUrl(task, options = {}) { if (!isKuaishouCloudTask(task)) { - throw createHttpError('当前任务不是快手 Cloud 履约任务', { + throw createHttpError("当前任务不是快手 Cloud 履约任务", { statusCode: 409, - errorCode: 'kuaishou_cloud_task_invalid', - }) + errorCode: "kuaishou_cloud_task_invalid", + }); } - const now = nowIso() - const actor = normalizeActor(options.actor) - const taskContext = parseTaskContext(task) - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) + const now = nowIso(); + const actor = normalizeActor(options.actor); + const taskContext = parseTaskContext(task); + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment); if (!flow.binding.vnId || !flow.binding.vnKey) { - throw createHttpError('当前任务缺少可刷新绑定链接的虚拟号信息', { + throw createHttpError("当前任务缺少可刷新绑定链接的虚拟号信息", { statusCode: 409, - errorCode: 'kuaishou_cloud_missing_bind_url_context', - }) + errorCode: "kuaishou_cloud_missing_bind_url_context", + }); } - const cloudContext = resolvePersistedCloudtentaclesContext() - const oldVnKey = flow.binding.vnKey - const oldVnId = flow.binding.vnId - const oldVnPhone = flow.binding.vnPhone + const cloudContext = resolvePersistedCloudtentaclesContext( + flow.binding.cloudSourceKey || "default" + ); + const oldVnKey = flow.binding.vnKey; + const oldVnId = flow.binding.vnId; + const oldVnPhone = flow.binding.vnPhone; await backCloudtentaclesVirtualNumber({ ...cloudContext, key: oldVnKey, id: oldVnId, - }) + }); - await createTaskEvent(task.id, 'kuaishou_cloud_expired_bind_number_returned', { - source: String(options.source || 'system').trim() || 'system', - vnKey: oldVnKey, - vnId: oldVnId, - vnPhoneMasked: maskPhone(oldVnPhone), - actor, - }, now) + await createTaskEvent( + task.id, + "kuaishou_cloud_expired_bind_number_returned", + { + source: String(options.source || "system").trim() || "system", + vnKey: oldVnKey, + vnId: oldVnId, + vnPhoneMasked: maskPhone(oldVnPhone), + actor, + }, + now + ); - let preparedBinding + let preparedBinding; try { preparedBinding = await prepareKuaishouCloudBindResourceWithFallback({ cloudContext, @@ -437,13 +522,13 @@ export async function refreshKuaishouCloudTaskBindUrl(task, options = {}) { flow, binding: flow.binding, }), - }) + }); - if (!String(preparedBinding.bindUrl || '').trim()) { - throw createHttpError('cloudtentacles 未返回新的绑定链接', { + if (!String(preparedBinding.bindUrl || "").trim()) { + throw createHttpError("cloudtentacles 未返回新的绑定链接", { statusCode: 502, - errorCode: 'kuaishou_cloud_empty_bind_url', - }) + errorCode: "kuaishou_cloud_empty_bind_url", + }); } } catch (error) { const failedTask = await markKuaishouCloudBindUrlRefreshFailed(task, { @@ -452,14 +537,18 @@ export async function refreshKuaishouCloudTaskBindUrl(task, options = {}) { now, actor, error, - }) + }); return { task: failedTask, - claimUrl: buildClaimUrl(String(task.primary_claim_token || task.claim_token || '')), - token: String(task.primary_claim_token || task.claim_token || ''), - flow: normalizeKuaishouCloudFlow(parseTaskContext(failedTask).kuaishouCloudFulfillment), - } + claimUrl: buildClaimUrl( + String(task.primary_claim_token || task.claim_token || "") + ), + token: String(task.primary_claim_token || task.claim_token || ""), + flow: normalizeKuaishouCloudFlow( + parseTaskContext(failedTask).kuaishouCloudFulfillment + ), + }; } const nextContext = { @@ -468,7 +557,7 @@ export async function refreshKuaishouCloudTaskBindUrl(task, options = {}) { ...flow, binding: { ...flow.binding, - prepareStatus: 'ready', + prepareStatus: "ready", vnKey: preparedBinding.vnKey, vnId: preparedBinding.vnId, vnPhone: preparedBinding.vnPhone, @@ -476,100 +565,128 @@ export async function refreshKuaishouCloudTaskBindUrl(task, options = {}) { bindPreparedAt: now, bindExpiresAt: resolveKuaishouCloudBindUrlExpiresAt(now), bindProbeAt: null, - bindProbeStatus: 'pending', - bindProbeMessage: '', - roleName: String(task.task_status || '').trim() === 'role_confirmed' ? flow.binding.roleName : '', - roleId: String(task.task_status || '').trim() === 'role_confirmed' ? flow.binding.roleId : '', + bindProbeStatus: "pending", + bindProbeMessage: "", + roleName: + String(task.task_status || "").trim() === "role_confirmed" + ? flow.binding.roleName + : "", + roleId: + String(task.task_status || "").trim() === "role_confirmed" + ? flow.binding.roleId + : "", }, - role: String(task.task_status || '').trim() === 'role_confirmed' - ? flow.role - : { - status: 'pending', - name: '', - rid: '', - refreshedAt: null, - errorMessage: '', - rawInfo: null, - }, + role: + String(task.task_status || "").trim() === "role_confirmed" + ? flow.role + : { + status: "pending", + name: "", + rid: "", + refreshedAt: null, + errorMessage: "", + rawInfo: null, + }, }, - } + }; - const claimLinkState = options.claimLinkState || await ensureTaskClaimLink(task) + const claimLinkState = + options.claimLinkState || (await ensureTaskClaimLink(task)); const updatedTask = await updateTask(task.id, { - task_status: String(task.task_status || '').trim() === 'role_confirmed' ? task.task_status : 'waiting_binding', - claim_token: claimLinkState.token || task.claim_token || '', + task_status: + String(task.task_status || "").trim() === "role_confirmed" + ? task.task_status + : "waiting_binding", + claim_token: claimLinkState.token || task.claim_token || "", claim_expires_at: claimLinkState.expiredAt || getTaskClaimExpiresAt(task), - role_id: String(task.task_status || '').trim() === 'role_confirmed' ? task.role_id : '', - role_name: String(task.task_status || '').trim() === 'role_confirmed' ? task.role_name : '', - last_error: '', + role_id: + String(task.task_status || "").trim() === "role_confirmed" + ? task.role_id + : "", + role_name: + String(task.task_status || "").trim() === "role_confirmed" + ? task.role_name + : "", + last_error: "", context_json: JSON.stringify(nextContext), updated_at: now, - }) + }); - await createTaskEvent(task.id, 'kuaishou_cloud_bind_url_refreshed', { - source: String(options.source || 'system').trim() || 'system', - oldVnId, - oldVnPhoneMasked: maskPhone(oldVnPhone), - vnKey: preparedBinding.vnKey, - vnId: preparedBinding.vnId, - vnPhoneMasked: maskPhone(preparedBinding.vnPhone), - actor, - }, now) + await createTaskEvent( + task.id, + "kuaishou_cloud_bind_url_refreshed", + { + source: String(options.source || "system").trim() || "system", + oldVnId, + oldVnPhoneMasked: maskPhone(oldVnPhone), + vnKey: preparedBinding.vnKey, + vnId: preparedBinding.vnId, + vnPhoneMasked: maskPhone(preparedBinding.vnPhone), + actor, + }, + now + ); return { task: updatedTask, claimUrl: claimLinkState.claimUrl, token: claimLinkState.token, flow: normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment), - } + }; } export async function probeKuaishouCloudTaskBindUrl(task, options = {}) { if (!isKuaishouCloudTask(task)) { - throw createHttpError('当前任务不是快手 Cloud 履约任务', { + throw createHttpError("当前任务不是快手 Cloud 履约任务", { statusCode: 409, - errorCode: 'kuaishou_cloud_task_invalid', - }) + errorCode: "kuaishou_cloud_task_invalid", + }); } - const now = nowIso() - const taskContext = parseTaskContext(task) - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) + const now = nowIso(); + const taskContext = parseTaskContext(task); + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment); if (!flow.binding.bindUrl) { return { task, flow, probe: null, - } + }; } - const probeIntervalMs = Number(resolveCloudtentaclesConfig().bindUrlProbeIntervalSeconds || 30) * 1000 - if (!options.force && flow.binding.bindProbeAt && Date.now() - Date.parse(flow.binding.bindProbeAt) < probeIntervalMs) { + const probeIntervalMs = + Number(resolveCloudtentaclesConfig().bindUrlProbeIntervalSeconds || 30) * + 1000; + if ( + !options.force && + flow.binding.bindProbeAt && + Date.now() - Date.parse(flow.binding.bindProbeAt) < probeIntervalMs + ) { return { task, flow, probe: null, - } + }; } const probe = await probeCloudtentaclesBindUrl({ bindUrl: flow.binding.bindUrl, - }) + }); if (probe.expired || !isKuaishouCloudBindUrlFresh(flow)) { const refreshed = await refreshKuaishouCloudTaskBindUrl(task, { - source: options.source || 'claim_page_bind_url_expired', - actor: options.actor || { source: 'system' }, - }) + source: options.source || "claim_page_bind_url_expired", + actor: options.actor || { source: "system" }, + }); return { ...refreshed, probe, - } + }; } - const probeRoleInfo = normalizeKuaishouCloudRoleInfo(probe.roleInfo) - const hasRoleInfo = Boolean(probeRoleInfo.name || probeRoleInfo.rid) + const probeRoleInfo = normalizeKuaishouCloudRoleInfo(probe.roleInfo); + const hasRoleInfo = Boolean(probeRoleInfo.name || probeRoleInfo.rid); const nextContext = { ...taskContext, kuaishouCloudFulfillment: { @@ -577,137 +694,151 @@ export async function probeKuaishouCloudTaskBindUrl(task, options = {}) { binding: { ...flow.binding, bindProbeAt: now, - bindProbeStatus: probe.valid ? 'valid' : 'invalid', - bindProbeMessage: String(probe.message || probe.reason || '').trim(), + bindProbeStatus: probe.valid ? "valid" : "invalid", + bindProbeMessage: String(probe.message || probe.reason || "").trim(), roleName: hasRoleInfo ? probeRoleInfo.name : flow.binding.roleName, roleId: hasRoleInfo ? probeRoleInfo.rid : flow.binding.roleId, }, role: { ...flow.role, - status: hasRoleInfo ? 'ready' : flow.role.status, + status: hasRoleInfo ? "ready" : flow.role.status, name: hasRoleInfo ? probeRoleInfo.name : flow.role.name, rid: hasRoleInfo ? probeRoleInfo.rid : flow.role.rid, refreshedAt: hasRoleInfo ? now : flow.role.refreshedAt, - errorMessage: hasRoleInfo ? '' : flow.role.errorMessage, + errorMessage: hasRoleInfo ? "" : flow.role.errorMessage, rawInfo: hasRoleInfo ? probeRoleInfo.rawInfo : flow.role.rawInfo, }, }, - } + }; const updatedTask = await updateTask(task.id, { role_id: hasRoleInfo ? probeRoleInfo.rid : task.role_id, role_name: hasRoleInfo ? probeRoleInfo.name : task.role_name, context_json: JSON.stringify(nextContext), updated_at: now, - }) + }); return { task: updatedTask, flow: normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment), probe, - } + }; } /** * @param {any} task * @param {{ taskContext?: any, flow?: any, now?: string, actor?: any, error?: unknown }} [input] */ -async function markKuaishouCloudBindUrlRefreshFailed(task, { - taskContext, - flow, - now, - actor, - error, -} = {}) { - const errorMessage = error instanceof Error ? error.message : String(error || '新绑定链接准备失败') +async function markKuaishouCloudBindUrlRefreshFailed( + task, + { taskContext, flow, now, actor, error } = {} +) { + const errorMessage = + error instanceof Error + ? error.message + : String(error || "新绑定链接准备失败"); const nextContext = { ...taskContext, kuaishouCloudFulfillment: { ...flow, binding: { ...flow.binding, - prepareStatus: 'pending', + prepareStatus: "pending", vnId: 0, - vnPhone: '', - bindUrl: '', + vnPhone: "", + bindUrl: "", bindPreparedAt: null, bindExpiresAt: null, bindProbeAt: now, - bindProbeStatus: 'refresh_failed', + bindProbeStatus: "refresh_failed", bindProbeMessage: errorMessage, - roleName: '', - roleId: '', + roleName: "", + roleId: "", }, role: { - status: 'pending', - name: '', - rid: '', + status: "pending", + name: "", + rid: "", refreshedAt: now, - errorMessage: '绑定链接已过期,旧号码已退还,新链接准备失败,请稍后刷新或联系客服处理', + errorMessage: + "绑定链接已过期,旧号码已退还,新链接准备失败,请稍后刷新或联系客服处理", rawInfo: null, }, }, - } + }; const updatedTask = await updateTask(task.id, { - task_status: 'pending_binding_prepare', - role_id: '', - role_name: '', + task_status: "pending_binding_prepare", + role_id: "", + role_name: "", last_error: `绑定链接过期,旧号码已退还,新链接准备失败:${errorMessage}`, context_json: JSON.stringify(nextContext), updated_at: now, - }) + }); - await createTaskEvent(task.id, 'kuaishou_cloud_bind_url_refresh_failed', { - source: 'system_refresh_expired_bind_url', - errorMessage, - actor, - }, now) + await createTaskEvent( + task.id, + "kuaishou_cloud_bind_url_refresh_failed", + { + source: "system_refresh_expired_bind_url", + errorMessage, + actor, + }, + now + ); await notifyKuaishouCloudBindUrlRefreshFailed({ task: updatedTask, errorMessage, - }) + }); - return updatedTask + return updatedTask; } export async function refreshKuaishouCloudTaskRoleInfo(task, options = {}) { if (!isKuaishouCloudTask(task)) { - throw createHttpError('当前任务不是快手 Cloud 履约任务', { + throw createHttpError("当前任务不是快手 Cloud 履约任务", { statusCode: 409, - errorCode: 'kuaishou_cloud_task_invalid', - }) + errorCode: "kuaishou_cloud_task_invalid", + }); } - const now = nowIso() - const actor = normalizeActor(options.actor) - const taskContext = parseTaskContext(task) - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) + const now = nowIso(); + const actor = normalizeActor(options.actor); + const taskContext = parseTaskContext(task); + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment); if (!flow.binding.vnId || !flow.binding.vnKey) { - throw createHttpError('当前任务还没有可查询的绑定角色信息,请先准备绑定资源', { - statusCode: 409, - errorCode: 'kuaishou_cloud_missing_bind_info_context', - }) + throw createHttpError( + "当前任务还没有可查询的绑定角色信息,请先准备绑定资源", + { + statusCode: 409, + errorCode: "kuaishou_cloud_missing_bind_info_context", + } + ); } if (flow.binding.bindUrl) { const probed = await probeKuaishouCloudTaskBindUrl(task, { - source: options.source || 'system_role_refresh_bind_url_probe', + source: options.source || "system_role_refresh_bind_url_probe", actor, force: options.forceProbe === true, - }) - const probedFlow = normalizeKuaishouCloudFlow(parseTaskContext(probed.task).kuaishouCloudFulfillment) + }); + const probedFlow = normalizeKuaishouCloudFlow( + parseTaskContext(probed.task).kuaishouCloudFulfillment + ); if (probed.probe?.expired) { return { task: probed.task, roleInfo: normalizeKuaishouCloudRoleInfo(null), flow: probedFlow, - } + }; } - if ((probed.probe === null || probed.probe?.valid) && (probedFlow.binding.roleName || probedFlow.binding.roleId)) { + if ( + (probed.probe === null || probed.probe?.valid) && + (probedFlow.binding.roleName || probedFlow.binding.roleId) + ) { return { task: probed.task, roleInfo: { @@ -716,95 +847,112 @@ export async function refreshKuaishouCloudTaskRoleInfo(task, options = {}) { rawInfo: probedFlow.role.rawInfo, }, flow: probedFlow, - } + }; } } - const cloudContext = resolvePersistedCloudtentaclesContext() + const cloudContext = resolvePersistedCloudtentaclesContext( + flow.binding.cloudSourceKey || "default" + ); const bindInfoResult = await getCloudtentaclesBindInfo({ ...cloudContext, key: flow.binding.vnKey, id: flow.binding.vnId, - }) + }); - const bindInfo = normalizeKuaishouCloudRoleInfo(bindInfoResult.bindInfo) - const hasRoleInfo = Boolean(bindInfo.name || bindInfo.rid) + const bindInfo = normalizeKuaishouCloudRoleInfo(bindInfoResult.bindInfo); + const hasRoleInfo = Boolean(bindInfo.name || bindInfo.rid); const nextContext = { ...taskContext, kuaishouCloudFulfillment: { ...flow, binding: { ...flow.binding, - roleName: hasRoleInfo ? bindInfo.name : '', - roleId: hasRoleInfo ? bindInfo.rid : '', + roleName: hasRoleInfo ? bindInfo.name : "", + roleId: hasRoleInfo ? bindInfo.rid : "", }, role: { - status: hasRoleInfo ? 'ready' : 'pending', - name: hasRoleInfo ? bindInfo.name : '', - rid: hasRoleInfo ? bindInfo.rid : '', + status: hasRoleInfo ? "ready" : "pending", + name: hasRoleInfo ? bindInfo.name : "", + rid: hasRoleInfo ? bindInfo.rid : "", refreshedAt: now, - errorMessage: hasRoleInfo ? '' : '当前还没有查询到角色信息,请完成绑定后稍等片刻再试', + errorMessage: hasRoleInfo + ? "" + : "当前还没有查询到角色信息,请完成绑定后稍等片刻再试", rawInfo: bindInfo.rawInfo, }, }, - } + }; const updatedTask = await updateTask(task.id, { - role_id: hasRoleInfo ? bindInfo.rid : '', - role_name: hasRoleInfo ? bindInfo.name : '', + role_id: hasRoleInfo ? bindInfo.rid : "", + role_name: hasRoleInfo ? bindInfo.name : "", context_json: JSON.stringify(nextContext), updated_at: now, - }) + }); if (options.recordEvent !== false) { - await createTaskEvent(task.id, 'kuaishou_cloud_role_info_refreshed', { - source: String(options.source || 'system').trim() || 'system', - roleName: bindInfo.name, - roleId: bindInfo.rid, - vnId: flow.binding.vnId, - actor, - }, now) + await createTaskEvent( + task.id, + "kuaishou_cloud_role_info_refreshed", + { + source: String(options.source || "system").trim() || "system", + roleName: bindInfo.name, + roleId: bindInfo.rid, + vnId: flow.binding.vnId, + actor, + }, + now + ); } return { task: updatedTask, roleInfo: bindInfo, flow: normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment), - } + }; } export async function dispatchKuaishouCloudFulfillmentTask(task, options = {}) { if (!isKuaishouCloudTask(task)) { - throw createHttpError('当前任务不是快手 Cloud 履约任务', { + throw createHttpError("当前任务不是快手 Cloud 履约任务", { statusCode: 409, - errorCode: 'kuaishou_cloud_task_invalid', - }) + errorCode: "kuaishou_cloud_task_invalid", + }); } - const actor = normalizeActor(options.actor) - const now = nowIso() - const taskContext = parseTaskContext(task) - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) - const cloudContext = resolvePersistedCloudtentaclesContext() + const actor = normalizeActor(options.actor); + const now = nowIso(); + const taskContext = parseTaskContext(task); + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment); + const cloudContext = resolvePersistedCloudtentaclesContext( + flow.binding.cloudSourceKey || "default" + ); if (!flow.binding.skuId || !flow.binding.vnId || !flow.binding.vnPhone) { - throw createHttpError('当前任务还没有准备好绑定资源,请先完成绑定资源准备', { - statusCode: 409, - errorCode: 'kuaishou_cloud_not_prepared', - }) + throw createHttpError( + "当前任务还没有准备好绑定资源,请先完成绑定资源准备", + { + statusCode: 409, + errorCode: "kuaishou_cloud_not_prepared", + } + ); } - const ticketCode = String(options.ticketCode || '').trim() - const persistedTicketCode = String(flow.ticket.code || '').trim() + const ticketCode = String(options.ticketCode || "").trim(); + const persistedTicketCode = String(flow.ticket.code || "").trim(); if (!persistedTicketCode && !ticketCode) { - throw createHttpError('客户还没有提交有效核销码,暂时不能继续兑换', { + throw createHttpError("客户还没有提交有效核销码,暂时不能继续兑换", { statusCode: 409, - errorCode: 'kuaishou_cloud_missing_ticket_code', - }) + errorCode: "kuaishou_cloud_missing_ticket_code", + }); } - if (flow.dispatch.status === 'success' && String(task.task_status || '').trim() === 'dispatched_pending_return') { - return { task, flow } + if ( + flow.dispatch.status === "success" && + String(task.task_status || "").trim() === "dispatched_pending_return" + ) { + return { task, flow }; } const dispatchResult = await useCloudtentaclesSku({ @@ -812,7 +960,7 @@ export async function dispatchKuaishouCloudFulfillmentTask(task, options = {}) { id: flow.binding.skuId, virtualNumberId: flow.binding.vnId, phone: flow.binding.vnPhone, - }) + }); const nextContext = { ...taskContext, @@ -826,138 +974,175 @@ export async function dispatchKuaishouCloudFulfillmentTask(task, options = {}) { }, dispatch: { ...flow.dispatch, - status: 'success', + status: "success", dispatchAt: now, dispatchBy: actor, sendType: Number(dispatchResult.sendType || 0) || 0, - note: String(dispatchResult.note || dispatchResult.responseMessage || 'cloudtentacles 发货成功').trim(), + note: String( + dispatchResult.note || + dispatchResult.responseMessage || + "cloudtentacles 发货成功" + ).trim(), }, }, - } + }; let 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(), - user_action_status: 'not_required', - last_error: '', + task_status: "dispatched_pending_return", + delivery_status: "delivered", + result_code: "kuaishou_cloud_dispatched", + result_message: String( + dispatchResult.responseMessage || + dispatchResult.note || + "cloudtentacles 发货成功" + ).trim(), + user_action_status: "not_required", + last_error: "", context_json: JSON.stringify(nextContext), updated_at: now, - }) + }); - await createTaskEvent(task.id, 'kuaishou_cloud_dispatched', { - source: String(options.source || 'system').trim() || 'system', - ticketCodeMasked: maskCode(ticketCode || persistedTicketCode), - skuId: flow.binding.skuId, - vnId: flow.binding.vnId, - vnPhoneMasked: maskPhone(flow.binding.vnPhone), - sendType: dispatchResult.sendType, - note: dispatchResult.note, - actor, - }, now) + await createTaskEvent( + task.id, + "kuaishou_cloud_dispatched", + { + source: String(options.source || "system").trim() || "system", + ticketCodeMasked: maskCode(ticketCode || persistedTicketCode), + skuId: flow.binding.skuId, + vnId: flow.binding.vnId, + vnPhoneMasked: maskPhone(flow.binding.vnPhone), + sendType: dispatchResult.sendType, + note: dispatchResult.note, + actor, + }, + now + ); - const shouldAutoFinalize = options.autoFinalize === true - && normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment).returnNumber.autoReturnEnabled === true + const shouldAutoFinalize = + options.autoFinalize === true && + normalizeKuaishouCloudFlow(nextContext.kuaishouCloudFulfillment) + .returnNumber.autoReturnEnabled === true; if (shouldAutoFinalize) { - const finalizeResult = await returnKuaishouCloudFulfillmentTask(updatedTask, { - actor, - source: options.source || 'system_auto_finalize', - }) - updatedTask = finalizeResult.task + const finalizeResult = await returnKuaishouCloudFulfillmentTask( + updatedTask, + { + actor, + source: options.source || "system_auto_finalize", + } + ); + updatedTask = finalizeResult.task; } return { task: updatedTask, - flow: normalizeKuaishouCloudFlow(parseTaskContext(updatedTask).kuaishouCloudFulfillment), - } + flow: normalizeKuaishouCloudFlow( + parseTaskContext(updatedTask).kuaishouCloudFulfillment + ), + }; } export async function returnKuaishouCloudFulfillmentTask(task, options = {}) { if (!isKuaishouCloudTask(task)) { - throw createHttpError('当前任务不是快手 Cloud 履约任务', { + throw createHttpError("当前任务不是快手 Cloud 履约任务", { statusCode: 409, - errorCode: 'kuaishou_cloud_task_invalid', - }) + errorCode: "kuaishou_cloud_task_invalid", + }); } - const actor = normalizeActor(options.actor) - const now = nowIso() - const taskContext = parseTaskContext(task) - const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment) - const cloudContext = resolvePersistedCloudtentaclesContext() + const actor = normalizeActor(options.actor); + const now = nowIso(); + const taskContext = parseTaskContext(task); + const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment); + const cloudContext = resolvePersistedCloudtentaclesContext( + flow.binding.cloudSourceKey || "default" + ); if (!flow.binding.vnId || !flow.binding.vnKey) { - throw createHttpError('当前任务缺少可退还的虚拟号信息', { + throw createHttpError("当前任务缺少可退还的虚拟号信息", { statusCode: 409, - errorCode: 'kuaishou_cloud_missing_return_context', - }) + errorCode: "kuaishou_cloud_missing_return_context", + }); } - if (flow.returnNumber.status === 'success' && ['completed', 'manual_review'].includes(String(task.task_status || '').trim())) { - return { task, flow } + if ( + flow.returnNumber.status === "success" && + ["completed", "manual_review"].includes( + String(task.task_status || "").trim() + ) + ) { + return { task, flow }; } await backCloudtentaclesVirtualNumber({ ...cloudContext, key: flow.binding.vnKey, id: flow.binding.vnId, - }) + }); - const order = await getOrderById(task.order_id) - const ticketCode = String(flow.ticket.code || '').trim() - const shopId = String(flow.consume.shopId || order?.shop_id || '').trim() - const shopName = String(flow.consume.shopName || order?.shop_name || '').trim() - const eticketSource = getKuaishouEticketSourceConfig() + const order = await getOrderById(task.order_id); + const ticketCode = String(flow.ticket.code || "").trim(); + const shopId = String(flow.consume.shopId || order?.shop_id || "").trim(); + const shopName = String( + flow.consume.shopName || order?.shop_name || "" + ).trim(); + const eticketSource = getKuaishouEticketSourceConfig(); const shopConfig = resolveKuaishouEticketShopConfig({ shopId, shopName, - }) + }); - let consumeStatus = 'pending' - let consumeErrorMessage = '' - let consumedAt = null - let nextTaskStatus = 'completed' - let nextResultCode = 'kuaishou_cloud_completed' - let nextResultMessage = 'cloudtentacles 发货、退号并完成快手核销' + let consumeStatus = "pending"; + let consumeErrorMessage = ""; + let consumedAt = null; + let nextTaskStatus = "completed"; + let nextResultCode = "kuaishou_cloud_completed"; + let nextResultMessage = "cloudtentacles 发货、退号并完成快手核销"; if (!order) { - consumeStatus = 'failed' - consumeErrorMessage = '任务关联订单不存在,无法执行快手核销' + consumeStatus = "failed"; + consumeErrorMessage = "任务关联订单不存在,无法执行快手核销"; } else if (!ticketCode) { - consumeStatus = 'failed' - consumeErrorMessage = '客户未提交有效核销码,无法执行快手核销' - } else if (!shopConfig || shopConfig.enabled === false || !String(shopConfig.cookie || '').trim()) { - consumeStatus = 'failed' - consumeErrorMessage = '订单对应快手小店缺少可用 Cookie,无法执行快手核销' + consumeStatus = "failed"; + consumeErrorMessage = "客户未提交有效核销码,无法执行快手核销"; + } else if ( + !shopConfig || + shopConfig.enabled === false || + !String(shopConfig.cookie || "").trim() + ) { + consumeStatus = "failed"; + consumeErrorMessage = "订单对应快手小店缺少可用 Cookie,无法执行快手核销"; } else { try { const consumeResult = await consumeKuaishouEticket({ baseUrl: eticketSource.baseUrl, cookie: shopConfig.cookie, eTicketId: ticketCode, - oid: String(flow.ticket.oid || '').trim(), - formToken: String(flow.ticket.formToken || '').trim(), - }) + oid: String(flow.ticket.oid || "").trim(), + formToken: String(flow.ticket.formToken || "").trim(), + }); if (consumeResult.consumed) { - consumeStatus = 'success' - consumedAt = now + consumeStatus = "success"; + consumedAt = now; } else { - consumeStatus = 'failed' - consumeErrorMessage = String(consumeResult.errorMessage || '快手核销失败').trim() + consumeStatus = "failed"; + consumeErrorMessage = String( + consumeResult.errorMessage || "快手核销失败" + ).trim(); } } catch (error) { - consumeStatus = 'failed' - consumeErrorMessage = error instanceof Error ? error.message : '快手核销失败' + consumeStatus = "failed"; + consumeErrorMessage = + error instanceof Error ? error.message : "快手核销失败"; } } - if (consumeStatus !== 'success') { - nextTaskStatus = 'manual_review' - nextResultCode = 'kuaishou_cloud_consume_failed' - nextResultMessage = consumeErrorMessage || '号码已退还,但快手核销未完成,请人工处理' + if (consumeStatus !== "success") { + nextTaskStatus = "manual_review"; + nextResultCode = "kuaishou_cloud_consume_failed"; + nextResultMessage = + consumeErrorMessage || "号码已退还,但快手核销未完成,请人工处理"; } const nextContext = { @@ -966,7 +1151,7 @@ export async function returnKuaishouCloudFulfillmentTask(task, options = {}) { ...flow, returnNumber: { ...flow.returnNumber, - status: 'success', + status: "success", returnedAt: now, returnedBy: actor, }, @@ -980,31 +1165,38 @@ export async function returnKuaishouCloudFulfillmentTask(task, options = {}) { errorMessage: consumeErrorMessage, }, }, - } + }; const updatedTask = await updateTask(task.id, { task_status: nextTaskStatus, - delivery_status: 'delivered', + delivery_status: "delivered", result_code: nextResultCode, result_message: nextResultMessage, - redeemed_at: consumeStatus === 'success' ? now : task.redeemed_at, + redeemed_at: consumeStatus === "success" ? now : task.redeemed_at, last_error: consumeErrorMessage, context_json: JSON.stringify(nextContext), updated_at: now, - }) - - await createTaskEvent(task.id, 'kuaishou_cloud_number_returned', { - source: String(options.source || 'system').trim() || 'system', - vnId: flow.binding.vnId, - vnPhoneMasked: maskPhone(flow.binding.vnPhone), - actor, - }, now) + }); await createTaskEvent( task.id, - consumeStatus === 'success' ? 'kuaishou_cloud_consumed' : 'kuaishou_cloud_consume_failed', + "kuaishou_cloud_number_returned", { - source: String(options.source || 'system').trim() || 'system', + source: String(options.source || "system").trim() || "system", + vnId: flow.binding.vnId, + vnPhoneMasked: maskPhone(flow.binding.vnPhone), + actor, + }, + now + ); + + await createTaskEvent( + task.id, + consumeStatus === "success" + ? "kuaishou_cloud_consumed" + : "kuaishou_cloud_consume_failed", + { + source: String(options.source || "system").trim() || "system", ticketCodeMasked: maskCode(ticketCode), shopId, shopName, @@ -1012,10 +1204,10 @@ export async function returnKuaishouCloudFulfillmentTask(task, options = {}) { errorMessage: consumeErrorMessage, actor, }, - now, - ) + now + ); - if (consumeStatus !== 'success') { + if (consumeStatus !== "success") { await notifyKuaishouCloudConsumeFailed({ task: updatedTask, order, @@ -1023,176 +1215,194 @@ export async function returnKuaishouCloudFulfillmentTask(task, options = {}) { shopId, shopName, errorMessage: consumeErrorMessage, - }) + }); } return { task: updatedTask, - flow: normalizeKuaishouCloudFlow(parseTaskContext(updatedTask).kuaishouCloudFulfillment), - } + flow: normalizeKuaishouCloudFlow( + parseTaskContext(updatedTask).kuaishouCloudFulfillment + ), + }; } export function maskPhone(value) { - const text = String(value || '').trim() + const text = String(value || "").trim(); if (!text) { - return '' + return ""; } if (text.length <= 7) { - return `${text.slice(0, 2)}***${text.slice(-2)}` + return `${text.slice(0, 2)}***${text.slice(-2)}`; } - return `${text.slice(0, 3)}****${text.slice(-4)}` + return `${text.slice(0, 3)}****${text.slice(-4)}`; } export function maskCode(value) { - const text = String(value || '').trim() + const text = String(value || "").trim(); if (!text) { - return '' + return ""; } if (text.length <= 8) { - return `${text.slice(0, 2)}***${text.slice(-2)}` + return `${text.slice(0, 2)}***${text.slice(-2)}`; } - return `${text.slice(0, 4)}****${text.slice(-4)}` + return `${text.slice(0, 4)}****${text.slice(-4)}`; } export function resolveKuaishouCloudBindUrlExpiresAt(preparedAt) { - const preparedTime = Date.parse(String(preparedAt || '')) + const preparedTime = Date.parse(String(preparedAt || "")); if (!Number.isFinite(preparedTime)) { - return null + return null; } - const config = resolveCloudtentaclesConfig() - const ttlSeconds = Number(config.bindUrlTtlSeconds || 600) - return new Date(preparedTime + Math.max(1, ttlSeconds) * 1000).toISOString() + const config = resolveCloudtentaclesConfig(); + const ttlSeconds = Number(config.bindUrlTtlSeconds || 600); + return new Date(preparedTime + Math.max(1, ttlSeconds) * 1000).toISOString(); } export function isKuaishouCloudBindUrlFresh(flow, now = new Date()) { - const normalizedFlow = normalizeKuaishouCloudFlow(flow) + const normalizedFlow = normalizeKuaishouCloudFlow(flow); if (!normalizedFlow.binding.bindUrl) { - return false + return false; } - const expiresAt = normalizedFlow.binding.bindExpiresAt + const expiresAt = normalizedFlow.binding.bindExpiresAt; if (!expiresAt) { - return false + return false; } - const expiresTime = Date.parse(String(expiresAt || '')) + const expiresTime = Date.parse(String(expiresAt || "")); if (!Number.isFinite(expiresTime)) { - return false + return false; } - return expiresTime > now.getTime() + return expiresTime > now.getTime(); } -function resolveKuaishouCloudBindingResources(flow, { skuItems = [], knapsackItems = [] } = {}) { - const normalizedSkuItems = Array.isArray(skuItems) ? skuItems.filter(isCloudSkuLikeItem) : [] - const normalizedKnapsackItems = Array.isArray(knapsackItems) ? knapsackItems.filter(isCloudSkuLikeItem) : [] - const currentSkuId = Number(flow?.binding?.skuId || 0) || 0 - const currentSkuName = String(flow?.binding?.skuName || '').trim() - const nameCandidates = collectKuaishouCloudNameCandidates(flow) +function resolveKuaishouCloudBindingResources( + flow, + { skuItems = [], knapsackItems = [] } = {} +) { + const normalizedSkuItems = Array.isArray(skuItems) + ? skuItems.filter(isCloudSkuLikeItem) + : []; + const normalizedKnapsackItems = Array.isArray(knapsackItems) + ? knapsackItems.filter(isCloudSkuLikeItem) + : []; + const currentSkuId = Number(flow?.binding?.skuId || 0) || 0; + const currentSkuName = String(flow?.binding?.skuName || "").trim(); + const nameCandidates = collectKuaishouCloudNameCandidates(flow); - const skuItemById = currentSkuId > 0 - ? normalizedSkuItems.find((item) => Number(item.id || 0) === currentSkuId) || null - : null - const knapsackItemById = currentSkuId > 0 - ? normalizedKnapsackItems.find((item) => Number(item.id || 0) === currentSkuId) || null - : null + 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 + const matchedItem = skuItemById || knapsackItemById; return { skuId: Number(matchedItem?.id || 0) || 0, - skuName: String(currentSkuName || matchedItem?.name || '').trim(), + skuName: String(currentSkuName || matchedItem?.name || "").trim(), vnKey: KUAISHOU_CLOUD_FIXED_VN_KEY, skuItem: skuItemById, knapsackItem: knapsackItemById, resolvedByName: false, - } + }; } - const matchedSkuItem = findCloudItemByNames(normalizedSkuItems, nameCandidates) + const matchedSkuItem = findCloudItemByNames( + normalizedSkuItems, + nameCandidates + ); const matchedKnapsackItem = findCloudItemByNames( normalizedKnapsackItems, nameCandidates, - matchedSkuItem ? Number(matchedSkuItem.id || 0) : 0, - ) - const matchedItem = matchedSkuItem || matchedKnapsackItem + matchedSkuItem ? Number(matchedSkuItem.id || 0) : 0 + ); + const matchedItem = matchedSkuItem || matchedKnapsackItem; return { skuId: Number(matchedItem?.id || 0) || 0, - skuName: String(currentSkuName || matchedItem?.name || '').trim(), + skuName: String(currentSkuName || matchedItem?.name || "").trim(), vnKey: KUAISHOU_CLOUD_FIXED_VN_KEY, skuItem: matchedSkuItem, knapsackItem: matchedKnapsackItem, resolvedByName: Boolean(matchedItem), - } + }; } function resolveKuaishouCloudVnKeyCandidates(_input = {}) { - return [KUAISHOU_CLOUD_FIXED_VN_KEY] + return [KUAISHOU_CLOUD_FIXED_VN_KEY]; } async function prepareKuaishouCloudBindResourceWithFallback(input = {}) { - const { cloudContext = {}, vnKeyCandidates = [] } = input - const candidates = Array.isArray(vnKeyCandidates) ? vnKeyCandidates : [] - let lastError = null + const { cloudContext = {}, vnKeyCandidates = [] } = input; + const candidates = Array.isArray(vnKeyCandidates) ? vnKeyCandidates : []; + let lastError = null; for (const vnKey of candidates) { - let vnId = 0 - let vnPhone = '' + 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() + }); + vnId = Number(appointed.item?.id || 0); + vnPhone = String(appointed.item?.phone || "").trim(); if (!vnId || !vnPhone) { - throw createHttpError('申请虚拟号成功但返回数据不完整', { + throw createHttpError("申请虚拟号成功但返回数据不完整", { statusCode: 502, - errorCode: 'kuaishou_cloud_invalid_vn', - }) + errorCode: "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(), - } + bindUrl: String(bindUrlResult.bindUrl || "").trim(), + }; } catch (error) { - lastError = error + lastError = error; if (vnId > 0) { try { @@ -1200,96 +1410,117 @@ async function prepareKuaishouCloudBindResourceWithFallback(input = {}) { ...cloudContext, key: vnKey, id: vnId, - }) + }); } catch { // 退号失败保留主错误 } } if (!isRecoverableKuaishouCloudVnKeyError(error)) { - throw error + throw error; } } } - throw lastError || createHttpError('没有找到可用的 VN Key', { - statusCode: 409, - errorCode: 'kuaishou_cloud_missing_binding_config', - }) + throw ( + lastError || + createHttpError("没有找到可用的 VN Key", { + statusCode: 409, + errorCode: "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))) + 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 normalizedItems = Array.isArray(items) ? items : []; const normalizedNames = nameCandidates .map((item) => ({ - raw: String(item || '').trim(), + raw: String(item || "").trim(), normalized: normalizeProductName(item), })) - .filter((item) => item.raw && item.normalized) + .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 + ? normalizedItems.find((item) => Number(item.id || 0) === preferredId) || + null + : null; } if (preferredId > 0) { - const preferred = normalizedItems.find((item) => Number(item.id || 0) === preferredId) || null + const preferred = + normalizedItems.find((item) => Number(item.id || 0) === preferredId) || + null; if (preferred) { - return preferred + return preferred; } } const exactMatches = normalizedItems.filter((item) => { - const itemName = normalizeProductName(item.name) - return normalizedNames.some((candidate) => candidate.normalized === itemName) - }) + const itemName = normalizeProductName(item.name); + return normalizedNames.some( + (candidate) => candidate.normalized === itemName + ); + }); if (exactMatches.length > 0) { - return exactMatches[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)) - }) + 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 partialMatches.sort( + (left, right) => + String(left.name || "").length - String(right.name || "").length + )[0]; } - return null + return null; } function isCloudSkuLikeItem(item) { - return Boolean(item) && typeof item === 'object' && Number(item.id || 0) > 0 + return Boolean(item) && typeof item === "object" && Number(item.id || 0) > 0; } function isRecoverableKuaishouCloudVnKeyError(error) { - const errorCode = String(error?.errorCode || error?.code || '').trim() - const errorMessage = String(error?.message || '').trim() - return errorCode === 'cloudtentacles_vn_bind_url_failed' - && errorMessage.includes('不支持的游戏类型') + const errorCode = String(error?.errorCode || error?.code || "").trim(); + const errorMessage = String(error?.message || "").trim(); + return ( + errorCode === "cloudtentacles_vn_bind_url_failed" && + errorMessage.includes("不支持的游戏类型") + ); } function normalizeActor(actor) { - if (!actor || typeof actor !== 'object') { - return null + if (!actor || typeof actor !== "object") { + return null; } - const source = String(actor.source || '').trim() - const userId = Number(actor.userId || 0) || 0 - const username = String(actor.username || '').trim() - const role = String(actor.role || '').trim() + const source = String(actor.source || "").trim(); + const userId = Number(actor.userId || 0) || 0; + const username = String(actor.username || "").trim(); + const role = String(actor.role || "").trim(); if (!source && !userId && !username && !role) { - return null + return null; } return { @@ -1297,36 +1528,36 @@ function normalizeActor(actor) { userId, username, role, - } + }; } function parseTaskContext(task) { - const value = task?.context_json + const value = task?.context_json; if (!value) { - return {} + return {}; } - if (typeof value === 'object') { - return value + if (typeof value === "object") { + return value; } try { - return JSON.parse(String(value || '{}')) + return JSON.parse(String(value || "{}")); } catch { - return {} + return {}; } } function getTaskClaimExpiresAt(task) { - return task?.claim_expires_at || task?.primary_claim_expires_at || null + return task?.claim_expires_at || task?.primary_claim_expires_at || null; } function isClaimExpired(expiredAt) { if (!expiredAt) { - return false + return false; } - const timestamp = new Date(expiredAt).getTime() - return Number.isFinite(timestamp) && timestamp <= Date.now() + const timestamp = new Date(expiredAt).getTime(); + return Number.isFinite(timestamp) && timestamp <= Date.now(); } diff --git a/apps/backend/src/services/platforms/cloudtentacles/session-state-service.js b/apps/backend/src/services/platforms/cloudtentacles/session-state-service.js index 5a4c96e9..4e7a447c 100644 --- a/apps/backend/src/services/platforms/cloudtentacles/session-state-service.js +++ b/apps/backend/src/services/platforms/cloudtentacles/session-state-service.js @@ -11,36 +11,101 @@ export function getCloudtentaclesSessionFilePath() { return CLOUDTENTACLES_SESSION_FILE_PATH } +/** + * Backward-compatible: returns the session for key='default'. + * Old callers that expect a single session object still work. + */ export function getCloudtentaclesSessionState() { - return loadCloudtentaclesSessionStateFromFile() + const states = loadCloudtentaclesSessionStatesFromFile() + return states.sessions['default'] || createDefaultCloudtentaclesSessionState() } +/** + * Get a session by its sourceKey. Returns null if not found. + */ +export function getCloudtentaclesSessionStateByKey(sourceKey) { + const states = loadCloudtentaclesSessionStatesFromFile() + const key = String(sourceKey || '').trim() + if (!key) return null + return states.sessions[key] || null +} + +/** + * Return the entire session states map: { sessions: { 'default': {...}, ... } }. + */ +export function getAllCloudtentaclesSessionStates() { + return loadCloudtentaclesSessionStatesFromFile() +} + +/** + * Backward-compatible save: accepts both old single-object format + * and new sessions-map format, normalizes, and persists. + */ export function saveCloudtentaclesSessionState(rawValue) { - const normalized = normalizeCloudtentaclesSessionState(rawValue) + const normalized = normalizeSessionStatesFile(rawValue) fs.mkdirSync(path.dirname(CLOUDTENTACLES_SESSION_FILE_PATH), { recursive: true }) fs.writeFileSync(CLOUDTENTACLES_SESSION_FILE_PATH, `${JSON.stringify(normalized, null, 2)}\n`, 'utf8') - return normalized + return normalized.sessions['default'] || createDefaultCloudtentaclesSessionState() } +/** + * Save a session for a specific sourceKey. + */ +export function saveCloudtentaclesSessionStateByKey(sourceKey, rawValue) { + const key = String(sourceKey || '').trim() + if (!key) { + throw new Error('saveCloudtentaclesSessionStateByKey: sourceKey is required') + } + + const states = loadCloudtentaclesSessionStatesFromFile() + states.sessions[key] = normalizeCloudtentaclesSessionState(rawValue) + + fs.mkdirSync(path.dirname(CLOUDTENTACLES_SESSION_FILE_PATH), { recursive: true }) + fs.writeFileSync(CLOUDTENTACLES_SESSION_FILE_PATH, `${JSON.stringify(states, null, 2)}\n`, 'utf8') + return states.sessions[key] +} + +/** + * Backward-compatible clear: clears the 'default' session only. + */ export function clearCloudtentaclesSessionState() { + return clearCloudtentaclesSessionStateByKey('default') +} + +/** + * Clear a session by sourceKey (sets it to default empty state). + */ +export function clearCloudtentaclesSessionStateByKey(sourceKey) { + const key = String(sourceKey || '').trim() + if (!key) { + throw new Error('clearCloudtentaclesSessionStateByKey: sourceKey is required') + } + const cleared = createDefaultCloudtentaclesSessionState() - saveCloudtentaclesSessionState(cleared) + saveCloudtentaclesSessionStateByKey(key, cleared) return cleared } -function loadCloudtentaclesSessionStateFromFile() { +// --------------------------------------------------------------------------- +// Internal helpers +// --------------------------------------------------------------------------- + +function loadCloudtentaclesSessionStatesFromFile() { if (!fs.existsSync(CLOUDTENTACLES_SESSION_FILE_PATH)) { - return createDefaultCloudtentaclesSessionState() + return createDefaultCloudtentaclesSessionStates() } try { const rawText = fs.readFileSync(CLOUDTENTACLES_SESSION_FILE_PATH, 'utf8') - return normalizeCloudtentaclesSessionState(JSON.parse(rawText)) + return normalizeSessionStatesFile(JSON.parse(rawText)) } catch { - return createDefaultCloudtentaclesSessionState() + return createDefaultCloudtentaclesSessionStates() } } +/** + * Normalize a single session state item. + */ function normalizeCloudtentaclesSessionState(rawValue) { const source = isPlainObject(rawValue) ? rawValue : {} @@ -55,6 +120,31 @@ function normalizeCloudtentaclesSessionState(rawValue) { } } +/** + * Normalize the overall session states file format. + * Handles old format (single object without sessions key) by auto-wrapping + * into { sessions: { 'default': ... } }. + */ +function normalizeSessionStatesFile(rawValue) { + // Old format: { token: 'xxx', ... } (single object, no sessions key) + if (isPlainObject(rawValue) && !rawValue.sessions) { + return { + sessions: { + 'default': normalizeCloudtentaclesSessionState(rawValue), + }, + } + } + + // New format: { sessions: { 'default': {...}, ... } } + return { + sessions: isPlainObject(rawValue?.sessions) + ? Object.fromEntries( + Object.entries(rawValue.sessions).map(([k, v]) => [k, normalizeCloudtentaclesSessionState(v)]) + ) + : {}, + } +} + function createDefaultCloudtentaclesSessionState() { return { token: '', @@ -67,6 +157,14 @@ function createDefaultCloudtentaclesSessionState() { } } +function createDefaultCloudtentaclesSessionStates() { + return { + sessions: { + 'default': createDefaultCloudtentaclesSessionState(), + }, + } +} + function normalizeInteger(value, fallback) { const parsed = Number(value) return Number.isInteger(parsed) ? parsed : fallback diff --git a/apps/backend/src/services/platforms/cloudtentacles/source-config-service.js b/apps/backend/src/services/platforms/cloudtentacles/source-config-service.js index d0b48016..5cce28fd 100644 --- a/apps/backend/src/services/platforms/cloudtentacles/source-config-service.js +++ b/apps/backend/src/services/platforms/cloudtentacles/source-config-service.js @@ -11,35 +11,128 @@ export function getCloudtentaclesSourcesFilePath() { return CLOUDTENTACLES_SOURCES_FILE_PATH } +/** + * Backward-compatible: returns the source with key='default'. + * Old callers that expect a single source object still work. + */ export function getCloudtentaclesSourceConfig() { - return loadCloudtentaclesSourceConfigFromFile() + const config = loadCloudtentaclesSourcesConfigFromFile() + const defaultSource = config.sources.find(s => s.key === 'default') + return defaultSource || normalizeCloudtentaclesSourceItem({ key: 'default' }) } +/** + * Get a source item by its key. Returns null if not found. + */ +export function getCloudtentaclesSourceByKey(sourceKey) { + const config = loadCloudtentaclesSourcesConfigFromFile() + const key = String(sourceKey || '').trim() + if (!key) return null + return config.sources.find(s => s.key === key) || null +} + +/** + * Return the entire normalized config: { enabled, sources }. + */ +export function listCloudtentaclesSources() { + return loadCloudtentaclesSourcesConfigFromFile() +} + +/** + * Backward-compatible save: accepts both old single-object format + * and new list format, normalizes, and persists. + */ export function saveCloudtentaclesSourceConfig(rawValue) { - const normalized = normalizeCloudtentaclesSourceConfig(rawValue) + const normalized = normalizeCloudtentaclesSourcesConfig(rawValue) fs.mkdirSync(path.dirname(CLOUDTENTACLES_SOURCES_FILE_PATH), { recursive: true }) fs.writeFileSync(CLOUDTENTACLES_SOURCES_FILE_PATH, `${JSON.stringify(normalized, null, 2)}\n`, 'utf8') return normalized } -function loadCloudtentaclesSourceConfigFromFile() { +/** + * Save the entire list-format config object: { enabled, sources: [...] }. + */ +export function saveCloudtentaclesSourcesList(rawValue) { + return saveCloudtentaclesSourceConfig(rawValue) +} + +/** + * Save or update a single source item identified by sourceKey. + * If a source with the same key exists, it is replaced; otherwise it is appended. + */ +export function saveCloudtentaclesSourceByKey(sourceKey, data) { + const key = String(sourceKey || '').trim() + if (!key) { + throw new Error('saveCloudtentaclesSourceByKey: sourceKey is required') + } + + const config = loadCloudtentaclesSourcesConfigFromFile() + const normalizedItem = normalizeCloudtentaclesSourceItem({ ...data, key }) + const existingIndex = config.sources.findIndex(s => s.key === key) + + if (existingIndex >= 0) { + config.sources[existingIndex] = normalizedItem + } else { + config.sources.push(normalizedItem) + } + + fs.mkdirSync(path.dirname(CLOUDTENTACLES_SOURCES_FILE_PATH), { recursive: true }) + fs.writeFileSync(CLOUDTENTACLES_SOURCES_FILE_PATH, `${JSON.stringify(config, null, 2)}\n`, 'utf8') + return config +} + +/** + * Delete a single source by key. Throws if key is 'default' (cannot delete default source). + */ +export function deleteCloudtentaclesSourceByKey(sourceKey) { + const key = String(sourceKey || '').trim() + if (!key) { + throw new Error('deleteCloudtentaclesSourceByKey: sourceKey is required') + } + if (key === 'default') { + throw new Error('deleteCloudtentaclesSourceByKey: cannot delete the default source') + } + + const config = loadCloudtentaclesSourcesConfigFromFile() + const existingIndex = config.sources.findIndex(s => s.key === key) + + if (existingIndex < 0) { + return config + } + + config.sources.splice(existingIndex, 1) + + fs.mkdirSync(path.dirname(CLOUDTENTACLES_SOURCES_FILE_PATH), { recursive: true }) + fs.writeFileSync(CLOUDTENTACLES_SOURCES_FILE_PATH, `${JSON.stringify(config, null, 2)}\n`, 'utf8') + return config +} + +// --------------------------------------------------------------------------- +// Internal helpers +// --------------------------------------------------------------------------- + +function loadCloudtentaclesSourcesConfigFromFile() { if (!fs.existsSync(CLOUDTENTACLES_SOURCES_FILE_PATH)) { - return createDefaultCloudtentaclesSourceConfig() + return createDefaultCloudtentaclesSourcesConfig() } try { const rawText = fs.readFileSync(CLOUDTENTACLES_SOURCES_FILE_PATH, 'utf8') - return normalizeCloudtentaclesSourceConfig(JSON.parse(rawText)) + return normalizeCloudtentaclesSourcesConfig(JSON.parse(rawText)) } catch { - return createDefaultCloudtentaclesSourceConfig() + return createDefaultCloudtentaclesSourcesConfig() } } -function normalizeCloudtentaclesSourceConfig(rawValue) { +/** + * Normalize a single source item. Adds key (required) and label (optional). + */ +function normalizeCloudtentaclesSourceItem(rawValue) { const source = isPlainObject(rawValue) ? rawValue : {} return { - enabled: typeof source.enabled === 'boolean' ? source.enabled : true, + key: String(source.key || 'default').trim() || 'default', + label: String(source.label || '').trim(), baseUrl: String(source.baseUrl || 'https://123.207.217.176').trim() || 'https://123.207.217.176', username: String(source.username || '').trim(), password: String(source.password || '').trim(), @@ -49,15 +142,40 @@ function normalizeCloudtentaclesSourceConfig(rawValue) { } } -function createDefaultCloudtentaclesSourceConfig() { +/** + * Normalize the overall config. Handles both old single-object format + * (auto-migrates to new list format) and new { enabled, sources } format. + */ +function normalizeCloudtentaclesSourcesConfig(rawValue) { + // Old format: { enabled: true, username: 'xxx', ... } (single object, no sources array) + if (isPlainObject(rawValue) && !Array.isArray(rawValue.sources)) { + return { + enabled: rawValue.enabled !== false, + sources: [ + normalizeCloudtentaclesSourceItem({ + key: 'default', + label: '默认账号', + ...rawValue, // old fields auto-map to default source + }), + ].filter(Boolean), + } + } + + // New format: { enabled, sources: [...] } + return { + enabled: isPlainObject(rawValue) ? rawValue.enabled !== false : true, + sources: isPlainObject(rawValue) && Array.isArray(rawValue.sources) + ? rawValue.sources.map(s => normalizeCloudtentaclesSourceItem(s)).filter(Boolean) + : [], + } +} + +function createDefaultCloudtentaclesSourcesConfig() { return { enabled: true, - baseUrl: 'https://123.207.217.176', - username: '', - password: '', - phone: '', - deviceId: '-', - deviceType: 0, + sources: [ + normalizeCloudtentaclesSourceItem({ key: 'default', label: '默认账号' }), + ], } }