Files
order_site/apps/backend/src/services/fulfillment/kuaishou-cloud/task-finalization.ts
T
2026-07-09 21:03:39 +08:00

941 lines
29 KiB
TypeScript

import { createTaskEvent } from "../../../repositories/task-event-repo.js";
import { updateTask } from "../../../repositories/task-repo.js";
import { createHttpError } from "../../../utils/http.js";
import { nowIso } from "../../../utils/time.js";
import { TASK_STATUS, normalizeTaskStatus, type TaskStatus } from "../../../domain/task-status.js";
import { notifyKuaishouCloudAssetNotEnough } from "../../notification/domain-notifications.js";
import {
buyCloudtentaclesSku,
getCloudtentaclesAsset,
listCloudtentaclesSku,
useCloudtentaclesSku,
} from "../../platforms/cloudtentacles/catalog-service.js";
import { getCloudtentaclesKnapsack } from "../../platforms/cloudtentacles/knapsack-service.js";
import { backCloudtentaclesVirtualNumber } from "../../platforms/cloudtentacles/virtual-number-service.js";
import { consumeKuaishouIndustryVouchersForTask } from "../../platforms/kuaishou-industry/voucher-service.js";
import {
isKuaishouCloudTask,
isIndustryEVoucherTask,
maskCode,
maskPhone,
normalizeKuaishouCloudFlow,
type JsonObject,
} from "./domain.js";
import { resolvePersistedCloudtentaclesContextBySourceKeys } from "./cloudtentacles-context.js";
import { syncKuaishouCloudRoleInfoBeforeDispatch } from "./dispatch-role-sync.js";
import { normalizeActor, parseTaskContext } from "./task-context.js";
import type { TaskRow } from "../../../types/repository/rows.js";
type DispatchDeliveryItem = {
cloudSkuId: number;
cloudSkuName: string;
quantity: number;
};
type DispatchResultItem = DispatchDeliveryItem & {
unitIndex: number;
sendType: number;
note: string;
responseMessage: string;
};
type DispatchStockItem = DispatchDeliveryItem & {
requiredCount: number;
knapsackCount: number;
purchasedCount: number;
};
type DispatchStockResult = {
usedKnapsack: boolean;
purchaseTriggered: boolean;
assetBefore: number;
assetAfter: number;
items: DispatchStockItem[];
};
const cloudSkuDispatchLocks = new Map<string, Promise<void>>();
export async function dispatchKuaishouCloudFulfillmentTask(
task: TaskRow,
options: JsonObject = {}
) {
if (!isKuaishouCloudTask(task)) {
throw createHttpError("当前任务不是 kuaishou-lewan 履约任务", {
statusCode: 409,
errorCode: "kuaishou_cloud_task_invalid",
});
}
const actor = normalizeActor(options.actor);
const now = nowIso();
const taskContext = parseTaskContext(task);
const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment);
const cloudContext = resolvePersistedCloudtentaclesContextBySourceKeys([
flow.binding.resolvedSourceKey,
...flow.binding.cloudSourceKeys,
]);
if (!flow.binding.skuId || !flow.binding.vnId || !flow.binding.vnPhone) {
throw createHttpError(
"当前任务还没有准备好绑定资源,请先完成绑定资源准备",
{
statusCode: 409,
errorCode: "kuaishou_cloud_not_prepared",
}
);
}
const persistedTicketCode = String(flow.ticket.code || "").trim();
const voucherContext = isPlainObject(taskContext.kuaishouIndustryVoucher)
? taskContext.kuaishouIndustryVoucher
: {};
const industryVoucherCode = String(
voucherContext.voucherCode || voucherContext.eticketId || ""
).trim();
const hasIndustryVoucherForDispatch =
Boolean(industryVoucherCode) || isIndustryEVoucherTask(task);
const resolvedTicketCode = persistedTicketCode || industryVoucherCode;
if (!resolvedTicketCode && !hasIndustryVoucherForDispatch) {
throw createHttpError("旧快手小店核销流程已停用,请改用行业电子凭证处理", {
statusCode: 409,
errorCode: "kuaishou_cloud_missing_ticket_code",
});
}
if (
flow.dispatch.status === "success" &&
normalizeTaskStatus(task.task_status) === TASK_STATUS.DISPATCHED_PENDING_RETURN
) {
return { task, flow };
}
const synced = await syncKuaishouCloudRoleInfoBeforeDispatch(task, {
now,
actor,
source: options.source || "system_before_dispatch",
cloudContext,
taskContext,
flow,
});
const syncedFlow = synced.flow;
const syncedTaskContext = synced.taskContext;
const deliveryItems = resolveDispatchDeliveryItems(syncedFlow);
if (deliveryItems.length === 0) {
throw createHttpError("当前任务缺少 cloud 发货物品配置", {
statusCode: 409,
errorCode: "kuaishou_cloud_missing_delivery_items",
});
}
let dispatchResults: DispatchResultItem[];
let stockResult: DispatchStockResult;
try {
({ dispatchResults, stockResult } = await withCloudSkuDispatchLocks(
resolveCloudSkuDispatchLockKeys(cloudContext.resolvedSourceKey, deliveryItems),
() => prepareStockAndDispatch({
task,
flow: syncedFlow,
cloudContext,
deliveryItems,
})
));
} catch (error) {
await markKuaishouCloudDispatchFailed(task, syncedTaskContext, syncedFlow, error, {
actor,
source: options.source || "system",
});
throw error;
}
const firstDispatchResult: DispatchResultItem = dispatchResults[0] || {
cloudSkuId: syncedFlow.binding.skuId,
cloudSkuName: syncedFlow.binding.skuName,
quantity: 1,
unitIndex: 1,
sendType: 0,
note: "",
responseMessage: "",
};
const dispatchSummary =
dispatchResults.length > 1
? `cloudtentacles 发货成功,共 ${dispatchResults.length} 次`
: String(
firstDispatchResult.responseMessage ||
firstDispatchResult.note ||
"cloudtentacles 发货成功"
).trim();
const nextContext = {
...syncedTaskContext,
kuaishouCloudFulfillment: {
...syncedFlow,
ticket: {
...syncedFlow.ticket,
code: resolvedTicketCode,
capturedAt: resolvedTicketCode
? syncedFlow.ticket.capturedAt || now
: syncedFlow.ticket.capturedAt,
capturedBy: syncedFlow.ticket.capturedBy,
},
dispatch: {
...syncedFlow.dispatch,
status: "success",
dispatchAt: now,
dispatchBy: actor,
sendType: Number(firstDispatchResult.sendType || 0) || 0,
note: dispatchSummary,
items: dispatchResults,
},
purchase: {
...syncedFlow.purchase,
usedKnapsack: stockResult.usedKnapsack,
purchaseTriggered: stockResult.purchaseTriggered,
assetBefore: stockResult.assetBefore,
assetAfter: stockResult.assetAfter,
purchaseAt: stockResult.purchaseTriggered
? now
: syncedFlow.purchase.purchaseAt,
items: stockResult.items,
},
},
};
let updatedTask = await updateTask(task.id, {
task_status: TASK_STATUS.DISPATCHED_PENDING_RETURN,
delivery_status: "delivered",
result_code: "kuaishou_cloud_dispatched",
result_message: dispatchSummary,
user_action_status: "not_required",
last_error: "",
context_json: JSON.stringify(nextContext),
updated_at: now,
});
if (!updatedTask) {
throw createHttpError("kuaishou-lewan 发货状态更新失败", {
statusCode: 500,
errorCode: "kuaishou_cloud_dispatch_update_failed",
});
}
await createTaskEvent(
task.id,
"kuaishou_cloud_dispatched",
{
source: String(options.source || "system").trim() || "system",
ticketCodeMasked: maskCode(resolvedTicketCode),
skuId: syncedFlow.binding.skuId,
vnId: syncedFlow.binding.vnId,
vnPhoneMasked: maskPhone(syncedFlow.binding.vnPhone),
sendType: firstDispatchResult.sendType,
note: dispatchSummary,
deliveryItems,
dispatchResults,
stockResult,
actor,
},
now
);
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;
}
return {
task: updatedTask,
flow: normalizeKuaishouCloudFlow(
parseTaskContext(updatedTask).kuaishouCloudFulfillment
),
};
}
async function markKuaishouCloudDispatchFailed(
task: TaskRow,
taskContext: JsonObject,
flow: JsonObject,
error: unknown,
options: JsonObject = {}
) {
const now = nowIso();
const errorMessage = resolveErrorMessage(error);
const errorCode = resolveErrorCode(error) || "kuaishou_cloud_dispatch_failed";
const failureContext = resolveDispatchFailureContext(error);
const stockResult = isPlainObject(failureContext.stockResult)
? (failureContext.stockResult as Partial<DispatchStockResult>)
: null;
const nextContext = {
...taskContext,
kuaishouCloudFulfillment: {
...flow,
dispatch: {
...flow.dispatch,
status: "failed",
failedAt: now,
failedStage: String(failureContext.stage || "").trim(),
errorCode,
errorMessage,
note: buildDispatchFailedMessage(errorMessage),
},
purchase: {
...flow.purchase,
...(stockResult
? {
usedKnapsack: stockResult.usedKnapsack,
purchaseTriggered: stockResult.purchaseTriggered,
assetBefore: stockResult.assetBefore,
assetAfter: stockResult.assetAfter,
items: stockResult.items,
}
: {}),
},
},
};
await updateTask(task.id, {
task_status: TASK_STATUS.MANUAL_REVIEW,
user_action_status: "not_required",
last_error: errorMessage,
result_code: errorCode,
result_message: buildDispatchFailedMessage(errorMessage),
context_json: JSON.stringify(nextContext),
updated_at: now,
});
await createTaskEvent(
task.id,
"kuaishou_cloud_dispatch_failed",
{
source: String(options.source || "system").trim() || "system",
actor: options.actor,
errorCode,
errorMessage,
failureContext,
},
now
);
}
async function prepareStockAndDispatch({
task,
flow,
cloudContext,
deliveryItems,
}: {
task: TaskRow;
flow: JsonObject;
cloudContext: JsonObject;
deliveryItems: DispatchDeliveryItem[];
}) {
const [knapsack, skuList] = await Promise.all([
getCloudtentaclesKnapsack(cloudContext),
listCloudtentaclesSku(cloudContext),
]);
const skuItems = Array.isArray(skuList.items) ? skuList.items : [];
const knapsackItems = Array.isArray(knapsack.items) ? knapsack.items : [];
const stockItems = buildDispatchStockItems(deliveryItems, {
skuItems,
knapsackItems,
});
const missingItems = stockItems.filter((item) => item.purchasedCount > 0);
let assetBefore = 0;
let assetAfter = 0;
let purchaseTriggered = false;
if (missingItems.length > 0) {
if (flow.purchase?.autoBuyEnabled === false) {
throw createHttpError("背包中没有现成库存,且当前配置未开启自动购买", {
statusCode: 409,
errorCode: "kuaishou_cloud_auto_buy_disabled",
});
}
const missingSku = missingItems.find((item) => !findCloudSkuItem(skuItems, item.cloudSkuId));
if (missingSku) {
throw createHttpError(`cloudtentacles 未找到 SKU ${missingSku.cloudSkuId}`, {
statusCode: 404,
errorCode: "kuaishou_cloud_sku_not_found",
});
}
const asset = await getCloudtentaclesAsset(cloudContext);
assetBefore = Number(asset.asset || 0) || 0;
const targetPrice = missingItems.reduce((sum, item) => {
const skuItem = findCloudSkuItem(skuItems, item.cloudSkuId);
return sum + item.purchasedCount * Number(skuItem?.price || 0);
}, 0);
const requiredAsset = targetPrice + Number(flow.purchase?.minAssetReserve || 0);
if (assetBefore < requiredAsset) {
await notifyKuaishouCloudAssetNotEnough({
task,
flow,
assetBefore,
requiredAsset,
skuName: missingItems
.map((item) => item.cloudSkuName || `SKU ${item.cloudSkuId}`)
.join("、"),
});
throw createHttpError(`余额不足,当前 ${assetBefore},至少需要 ${requiredAsset}`, {
statusCode: 409,
errorCode: "kuaishou_cloud_asset_not_enough",
});
}
for (const item of missingItems) {
await createTaskEvent(
task.id,
"cloudtentacles_sku_buy_started",
{
source: "kuaishou_cloud_dispatch",
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
count: item.purchasedCount,
assetBefore,
},
nowIso()
);
try {
const buyResult = await buyCloudtentaclesSku({
...cloudContext,
id: item.cloudSkuId,
count: item.purchasedCount,
});
await createTaskEvent(
task.id,
"cloudtentacles_sku_buy_succeeded",
{
source: "kuaishou_cloud_dispatch",
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
count: item.purchasedCount,
responseMessage: buyResult.responseMessage,
},
nowIso()
);
} catch (error) {
await createTaskEvent(
task.id,
"cloudtentacles_sku_buy_failed",
{
source: "kuaishou_cloud_dispatch",
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
count: item.purchasedCount,
errorCode: resolveErrorCode(error),
errorMessage: resolveErrorMessage(error),
},
nowIso()
);
throw enrichCloudtentaclesDispatchError(error, {
stage: "purchase",
stockResult: buildCurrentStockResult({
stockItems,
purchaseTriggered,
assetBefore,
assetAfter,
}),
item,
});
}
}
purchaseTriggered = true;
const assetResult = await getCloudtentaclesAsset(cloudContext);
assetAfter = Number(assetResult.asset || 0) || 0;
}
const dispatchResults: DispatchResultItem[] = [];
for (const item of deliveryItems) {
for (let index = 0; index < item.quantity; index += 1) {
await createTaskEvent(
task.id,
"cloudtentacles_sku_use_started",
{
source: "kuaishou_cloud_dispatch",
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
unitIndex: index + 1,
quantity: item.quantity,
vnId: flow.binding.vnId,
vnPhoneMasked: maskPhone(flow.binding.vnPhone),
},
nowIso()
);
let dispatchResult: Awaited<ReturnType<typeof useCloudtentaclesSku>>;
try {
dispatchResult = await useCloudtentaclesSku({
...cloudContext,
id: item.cloudSkuId,
virtualNumberId: flow.binding.vnId,
phone: flow.binding.vnPhone,
});
} catch (error) {
await createTaskEvent(
task.id,
"cloudtentacles_sku_use_failed",
{
source: "kuaishou_cloud_dispatch",
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
unitIndex: index + 1,
quantity: item.quantity,
vnId: flow.binding.vnId,
vnPhoneMasked: maskPhone(flow.binding.vnPhone),
errorCode: resolveErrorCode(error),
errorMessage: resolveErrorMessage(error),
},
nowIso()
);
throw enrichCloudtentaclesDispatchError(error, {
stage: "dispatch",
stockResult: buildCurrentStockResult({
stockItems,
purchaseTriggered,
assetBefore,
assetAfter,
}),
item,
unitIndex: index + 1,
vnId: flow.binding.vnId,
vnPhoneMasked: maskPhone(flow.binding.vnPhone),
});
}
await createTaskEvent(
task.id,
"cloudtentacles_sku_use_succeeded",
{
source: "kuaishou_cloud_dispatch",
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
unitIndex: index + 1,
quantity: item.quantity,
vnId: flow.binding.vnId,
vnPhoneMasked: maskPhone(flow.binding.vnPhone),
sendType: Number(dispatchResult.sendType || 0) || 0,
note: String(dispatchResult.note || "").trim(),
responseMessage: String(dispatchResult.responseMessage || "").trim(),
},
nowIso()
);
dispatchResults.push({
cloudSkuId: item.cloudSkuId,
cloudSkuName: item.cloudSkuName,
unitIndex: index + 1,
quantity: item.quantity,
sendType: Number(dispatchResult.sendType || 0) || 0,
note: String(dispatchResult.note || "").trim(),
responseMessage: String(dispatchResult.responseMessage || "").trim(),
});
}
}
return {
dispatchResults,
stockResult: {
usedKnapsack: stockItems.every((item) => item.purchasedCount <= 0),
purchaseTriggered,
assetBefore,
assetAfter,
items: stockItems,
},
};
}
function buildCurrentStockResult({
stockItems,
purchaseTriggered,
assetBefore,
assetAfter,
}: {
stockItems: DispatchStockItem[];
purchaseTriggered: boolean;
assetBefore: number;
assetAfter: number;
}): DispatchStockResult {
return {
usedKnapsack: stockItems.every((item) => item.purchasedCount <= 0),
purchaseTriggered,
assetBefore,
assetAfter,
items: stockItems,
};
}
function enrichCloudtentaclesDispatchError(error: unknown, context: JsonObject) {
const currentError = error instanceof Error
? (error as Error & { context?: unknown })
: createHttpError(resolveErrorMessage(error), {
statusCode: 500,
errorCode: "kuaishou_cloud_dispatch_failed",
}) as Error & { context?: unknown };
const currentContext = isPlainObject(currentError.context) ? currentError.context : {};
currentError.context = {
...currentContext,
dispatchFailure: {
...(isPlainObject((currentContext as JsonObject).dispatchFailure)
? (currentContext as JsonObject).dispatchFailure
: {}),
...context,
},
};
return currentError;
}
function resolveDispatchFailureContext(error: unknown): JsonObject {
if (!error || typeof error !== "object") {
return {};
}
const context = (error as { context?: unknown }).context;
if (!isPlainObject(context)) {
return {};
}
const dispatchFailure = context.dispatchFailure;
return isPlainObject(dispatchFailure) ? dispatchFailure : {};
}
function resolveErrorCode(error: unknown) {
if (!error || typeof error !== "object") {
return "";
}
return String((error as { errorCode?: unknown }).errorCode || "").trim();
}
function resolveErrorMessage(error: unknown) {
return error instanceof Error ? error.message : String(error || "CloudTentacles 发货失败");
}
function buildDispatchFailedMessage(message: unknown) {
const reason = String(message || "").trim() || "未知错误";
return `CloudTentacles 发货失败:${reason}`;
}
function isPlainObject(value: unknown): value is JsonObject {
return Boolean(value) && typeof value === "object" && !Array.isArray(value);
}
function hasKuaishouIndustryVoucherContext(value: JsonObject): boolean {
const voucher = isPlainObject(value.kuaishouIndustryVoucher)
? value.kuaishouIndustryVoucher
: {};
const voucherCode = String(voucher.voucherCode || voucher.eticketId || "").trim();
const status = String(voucher.status || "UNUSED").trim().toUpperCase();
const sendCallbackStatus = String(voucher.sendCallbackStatus || "success").trim().toLowerCase();
return Boolean(voucherCode && status !== "DESTROYED" && sendCallbackStatus === "success");
}
export function buildDispatchStockItems(
deliveryItems: DispatchDeliveryItem[],
{ skuItems = [], knapsackItems = [] }: { skuItems?: JsonObject[]; knapsackItems?: JsonObject[] } = {}
): DispatchStockItem[] {
return deliveryItems.map((item) => {
const knapsackItem = findCloudSkuItem(knapsackItems, item.cloudSkuId);
const skuItem = findCloudSkuItem(skuItems, item.cloudSkuId);
const knapsackCount = Math.max(0, Number(knapsackItem?.count || 0) || 0);
return {
...item,
cloudSkuName: item.cloudSkuName || String(skuItem?.name || knapsackItem?.name || "").trim(),
requiredCount: item.quantity,
knapsackCount,
purchasedCount: Math.max(0, item.quantity - knapsackCount),
};
});
}
function findCloudSkuItem(items: JsonObject[], cloudSkuId: number) {
return items.find((item) => Number(item.id || 0) === cloudSkuId) || null;
}
function resolveCloudSkuDispatchLockKeys(sourceKey: unknown, deliveryItems: DispatchDeliveryItem[]) {
const normalizedSourceKey = String(sourceKey || "default").trim() || "default";
return deliveryItems.map((item) => `${normalizedSourceKey}:${item.cloudSkuId}`);
}
async function withCloudSkuDispatchLocks<T>(keys: string[], callback: () => Promise<T>): Promise<T> {
const normalizedKeys = [...new Set(keys.map((key) => String(key || "").trim()).filter(Boolean))].sort();
async function run(index: number): Promise<T> {
const key = normalizedKeys[index];
if (!key) {
return callback();
}
return withCloudSkuDispatchLock(key, () => run(index + 1));
}
return run(0);
}
async function withCloudSkuDispatchLock<T>(key: string, callback: () => Promise<T>): Promise<T> {
const previous = cloudSkuDispatchLocks.get(key) || Promise.resolve();
let release: () => void = () => {};
const current = new Promise<void>((resolve) => {
release = resolve;
});
const queued = previous.catch(() => {}).then(() => current);
cloudSkuDispatchLocks.set(key, queued);
await previous.catch(() => {});
try {
return await callback();
} finally {
release();
if (cloudSkuDispatchLocks.get(key) === queued) {
cloudSkuDispatchLocks.delete(key);
}
}
}
function resolveDispatchDeliveryItems(flow: JsonObject): DispatchDeliveryItem[] {
const rawItems = Array.isArray(flow.deliveryItems) ? flow.deliveryItems : [];
const items = rawItems
.map((item: unknown) => {
const source = item && typeof item === "object" ? item as JsonObject : {};
const cloudSkuId = Number(source.cloudSkuId || source.skuId || 0) || 0;
const quantity = Number(source.quantity || 1) || 1;
if (!Number.isInteger(cloudSkuId) || cloudSkuId <= 0) {
return null;
}
return {
cloudSkuId,
cloudSkuName: String(source.cloudSkuName || source.skuName || "").trim(),
quantity: Number.isInteger(quantity) && quantity > 0 ? quantity : 1,
};
})
.filter((item): item is DispatchDeliveryItem => Boolean(item));
if (items.length > 0) {
return items;
}
const fallbackSkuId = Number(flow?.binding?.skuId || 0) || 0;
if (!fallbackSkuId) {
return [];
}
return [
{
cloudSkuId: fallbackSkuId,
cloudSkuName: String(flow?.binding?.skuName || "").trim(),
quantity: 1,
},
];
}
export async function returnKuaishouCloudFulfillmentTask(
task: TaskRow,
options: JsonObject = {}
) {
if (!isKuaishouCloudTask(task)) {
throw createHttpError("当前任务不是 kuaishou-lewan 履约任务", {
statusCode: 409,
errorCode: "kuaishou_cloud_task_invalid",
});
}
const actor = normalizeActor(options.actor);
const now = nowIso();
const taskContext = parseTaskContext(task);
const flow = normalizeKuaishouCloudFlow(taskContext.kuaishouCloudFulfillment);
const cloudContext = resolvePersistedCloudtentaclesContextBySourceKeys([
flow.binding.resolvedSourceKey,
...flow.binding.cloudSourceKeys,
]);
if (!flow.binding.vnId || !flow.binding.vnKey) {
throw createHttpError("当前任务缺少可退还的虚拟号信息", {
statusCode: 409,
errorCode: "kuaishou_cloud_missing_return_context",
});
}
const taskStatus = normalizeTaskStatus(task.task_status);
if (
flow.returnNumber.status === "success" &&
(taskStatus === TASK_STATUS.COMPLETED || taskStatus === TASK_STATUS.MANUAL_REVIEW)
) {
return { task, flow };
}
await backCloudtentaclesVirtualNumber({
...cloudContext,
key: flow.binding.vnKey,
id: flow.binding.vnId,
});
const ticketCode = String(flow.ticket.code || "").trim();
let consumeStatus = "pending";
let consumeErrorMessage = "";
let consumedAt = null;
let nextTaskStatus: TaskStatus = TASK_STATUS.COMPLETED;
let nextResultCode = "kuaishou_cloud_completed";
const consumeAlreadyCompleted = flow.consume.status === "success";
const hasIndustryVoucher = hasKuaishouIndustryVoucherContext(taskContext);
const isIndustryTask = isIndustryEVoucherTask(task);
const shouldConsumeIndustryVoucher = hasIndustryVoucher || isIndustryTask;
let nextResultMessage = shouldConsumeIndustryVoucher
? "cloudtentacles 发货、退号并完成电子凭证核销"
: "cloudtentacles 发货、退号并收口";
let industryVoucherContextPatch: JsonObject | null = null;
if (consumeAlreadyCompleted) {
consumeStatus = "success";
consumedAt = flow.consume.consumedAt || now;
} else if (shouldConsumeIndustryVoucher) {
const industryResult = hasIndustryVoucher
? await consumeKuaishouIndustryVouchersForTask(task, {
source: String(options.source || "system_auto_finalize").trim() || "system_auto_finalize",
token: String(taskContext.kuaishouIndustryVoucher?.token || "").trim(),
consumeTime: Date.now(),
})
: { ok: true, consumed: [], failed: [] };
if (industryResult.ok) {
consumeStatus = "success";
consumedAt = now;
const consumedVoucher = industryResult.consumed[0] || null;
if (consumedVoucher) {
industryVoucherContextPatch = {
oid: consumedVoucher.oid,
token: consumedVoucher.token,
eticketId: consumedVoucher.voucher_code,
voucherCode: consumedVoucher.voucher_code,
unitIndex: Number(consumedVoucher.unit_index || 0) || 0,
status: "CONSUMED",
validStartTime: Number(consumedVoucher.valid_start_time || 0) || 0,
validEndTime: Number(consumedVoucher.valid_end_time || 0) || 0,
consumedAt: now,
consumeSerialNum: consumedVoucher.consume_serial_num || `CONSUME-${consumedVoucher.voucher_code}`,
};
}
} else {
consumeStatus = "failed";
consumeErrorMessage =
industryResult.failed[0]?.errorMessage || "电子凭证核销回调失败,请人工处理";
}
} else {
consumeStatus = "skipped";
consumeErrorMessage = "";
}
if (consumeStatus !== "success") {
if (shouldConsumeIndustryVoucher) {
nextTaskStatus = TASK_STATUS.MANUAL_REVIEW;
nextResultCode = "kuaishou_cloud_consume_failed";
nextResultMessage =
consumeErrorMessage || "号码已退还,但电子凭证核销未完成,请人工处理";
} else {
nextResultCode = "kuaishou_cloud_completed_without_eticket_consume";
}
}
const nextContext = {
...taskContext,
...(industryVoucherContextPatch
? {
kuaishouIndustryVoucher: {
...(isPlainObject(taskContext.kuaishouIndustryVoucher)
? taskContext.kuaishouIndustryVoucher
: {}),
...industryVoucherContextPatch,
},
}
: {}),
kuaishouCloudFulfillment: {
...flow,
returnNumber: {
...flow.returnNumber,
status: "success",
returnedAt: now,
returnedBy: actor,
},
consume: {
...flow.consume,
status: consumeStatus,
shopId: flow.consume.shopId,
shopName: flow.consume.shopName,
autoConsumeEnabled: flow.consume.autoConsumeEnabled === true,
consumedAt,
errorMessage: consumeErrorMessage,
},
},
};
const updatedTask = await updateTask(task.id, {
task_status: nextTaskStatus,
delivery_status: "delivered",
result_code: nextResultCode,
result_message: nextResultMessage,
redeemed_at: consumeStatus === "failed" ? task.redeemed_at : now,
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
);
const consumeEventType = consumeAlreadyCompleted
? "kuaishou_cloud_consume_already_completed"
: consumeStatus === "success"
? "kuaishou_cloud_consumed"
: consumeStatus === "skipped"
? "kuaishou_cloud_consume_skipped"
: "kuaishou_cloud_consume_failed";
await createTaskEvent(
task.id,
consumeEventType,
{
source: String(options.source || "system").trim() || "system",
ticketCodeMasked: maskCode(ticketCode),
shopId: flow.consume.shopId,
shopName: flow.consume.shopName,
consumeStatus,
consumeMode: shouldConsumeIndustryVoucher ? "industry_voucher" : "legacy_writeoff_disabled",
errorMessage: consumeErrorMessage,
actor,
},
now
);
return {
task: updatedTask,
flow: normalizeKuaishouCloudFlow(
parseTaskContext(updatedTask).kuaishouCloudFulfillment
),
};
}