优化任务状态投影和部署文档

This commit is contained in:
yml
2026-05-28 11:12:40 +08:00
parent 9709643583
commit ccc350cfea
15 changed files with 759 additions and 55 deletions
@@ -75,6 +75,8 @@ CREATE TABLE IF NOT EXISTS orders (
CREATE INDEX IF NOT EXISTS idx_orders_platform_order_lookup
ON orders(provider, platform, platform_order_id, id DESC);
CREATE INDEX IF NOT EXISTS idx_orders_created_at_desc
ON orders(created_at DESC);
CREATE TABLE IF NOT EXISTS order_items (
id BIGSERIAL PRIMARY KEY,
@@ -127,6 +129,10 @@ CREATE TABLE IF NOT EXISTS fulfillment_tasks (
CREATE INDEX IF NOT EXISTS idx_fulfillment_tasks_status ON fulfillment_tasks(task_status, delivery_status);
CREATE INDEX IF NOT EXISTS idx_fulfillment_tasks_order ON fulfillment_tasks(order_id, order_item_id);
CREATE INDEX IF NOT EXISTS idx_fulfillment_tasks_claim_token ON fulfillment_tasks(claim_token);
CREATE INDEX IF NOT EXISTS idx_fulfillment_tasks_platform_order_lookup
ON fulfillment_tasks(provider, platform, platform_order_id, id DESC);
CREATE INDEX IF NOT EXISTS idx_fulfillment_tasks_created_at_desc
ON fulfillment_tasks(created_at DESC);
CREATE TABLE IF NOT EXISTS claim_tokens (
id BIGSERIAL PRIMARY KEY,
@@ -169,3 +175,28 @@ CREATE TABLE IF NOT EXISTS task_runtime_contexts (
created_at TIMESTAMPTZ NOT NULL,
updated_at TIMESTAMPTZ NOT NULL
);
CREATE TABLE IF NOT EXISTS kuaishou_cloud_task_states (
task_id BIGINT PRIMARY KEY REFERENCES fulfillment_tasks(id) ON DELETE CASCADE,
ticket_status TEXT NOT NULL DEFAULT 'pending',
ticket_code_masked TEXT NOT NULL DEFAULT '',
bind_url TEXT NOT NULL DEFAULT '',
vn_phone TEXT NOT NULL DEFAULT '',
role_name TEXT NOT NULL DEFAULT '',
role_id TEXT NOT NULL DEFAULT '',
dispatch_status TEXT NOT NULL DEFAULT 'pending',
consume_status TEXT NOT NULL DEFAULT 'pending',
source_key TEXT NOT NULL DEFAULT '',
updated_at TIMESTAMPTZ NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_kuaishou_cloud_task_states_role_id
ON kuaishou_cloud_task_states(role_id);
CREATE INDEX IF NOT EXISTS idx_kuaishou_cloud_task_states_source_key
ON kuaishou_cloud_task_states(source_key);
CREATE INDEX IF NOT EXISTS idx_kuaishou_cloud_task_states_dispatch_status
ON kuaishou_cloud_task_states(dispatch_status);
CREATE INDEX IF NOT EXISTS idx_kuaishou_cloud_task_states_consume_status
ON kuaishou_cloud_task_states(consume_status);
CREATE INDEX IF NOT EXISTS idx_kuaishou_cloud_task_states_updated_at_desc
ON kuaishou_cloud_task_states(updated_at DESC);
@@ -0,0 +1,111 @@
import { query } from '../db/client.js'
import { maskCode } from '../utils/masking.js'
import { parseTaskContext } from '../utils/task-json.js'
type QueryExecutor = (text: string, params?: unknown[]) => Promise<{ rows: unknown[] }>
type JsonRecord = Record<string, any>
export type KuaishouCloudTaskStateSyncInput = {
taskId: number | string
executorKey?: unknown
contextJson?: unknown
updatedAt?: unknown
}
export async function syncKuaishouCloudTaskStateForTask(
input: KuaishouCloudTaskStateSyncInput,
executor: QueryExecutor = query,
): Promise<void> {
const taskId = Number(input.taskId || 0)
if (!Number.isFinite(taskId) || taskId <= 0) {
return
}
if (String(input.executorKey || '').trim() !== 'kuaishou_ct_assisted') {
await executor('DELETE FROM kuaishou_cloud_task_states WHERE task_id = $1', [taskId])
return
}
const flow = resolveKuaishouCloudFlow(input.contextJson)
if (!flow) {
await executor('DELETE FROM kuaishou_cloud_task_states WHERE task_id = $1', [taskId])
return
}
const ticket = asRecord(flow.ticket)
const binding = asRecord(flow.binding)
const role = asRecord(flow.role)
const dispatch = asRecord(flow.dispatch)
const consume = asRecord(flow.consume)
const roleName = pickFirstNonEmpty([role.name, binding.roleName])
const roleId = pickFirstNonEmpty([role.rid, binding.roleId])
const updatedAt = String(input.updatedAt || '').trim() || new Date().toISOString()
await executor(
`
INSERT INTO kuaishou_cloud_task_states (
task_id,
ticket_status,
ticket_code_masked,
bind_url,
vn_phone,
role_name,
role_id,
dispatch_status,
consume_status,
source_key,
updated_at
) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)
ON CONFLICT (task_id) DO UPDATE
SET
ticket_status = EXCLUDED.ticket_status,
ticket_code_masked = EXCLUDED.ticket_code_masked,
bind_url = EXCLUDED.bind_url,
vn_phone = EXCLUDED.vn_phone,
role_name = EXCLUDED.role_name,
role_id = EXCLUDED.role_id,
dispatch_status = EXCLUDED.dispatch_status,
consume_status = EXCLUDED.consume_status,
source_key = EXCLUDED.source_key,
updated_at = EXCLUDED.updated_at
`,
[
taskId,
normalizeStatus(ticket.status),
maskCode(ticket.code),
String(binding.bindUrl || '').trim(),
String(binding.vnPhone || '').trim(),
roleName,
roleId,
normalizeStatus(dispatch.status),
normalizeStatus(consume.status),
String(binding.resolvedSourceKey || '').trim(),
updatedAt,
],
)
}
function resolveKuaishouCloudFlow(contextJson: unknown): JsonRecord | null {
const context = parseTaskContext({ context_json: contextJson })
const flow = asRecord(context.kuaishouCloudFulfillment)
return Object.keys(flow).length > 0 ? flow : null
}
function asRecord(value: unknown): JsonRecord {
return value && typeof value === 'object' && !Array.isArray(value) ? value as JsonRecord : {}
}
function normalizeStatus(value: unknown): string {
return String(value || 'pending').trim() || 'pending'
}
function pickFirstNonEmpty(values: unknown[]): string {
for (const value of values) {
const normalized = String(value || '').trim()
if (normalized) {
return normalized
}
}
return ''
}
+83 -41
View File
@@ -1,6 +1,7 @@
import { query, withTransaction } from '../db/client.js'
import type { PoolClient } from 'pg'
import { syncKuaishouCloudTaskStateForTask } from './kuaishou-cloud-task-state-repo.js'
import type {
TaskCreateInput,
TaskListQueryInput,
@@ -42,13 +43,24 @@ const TASK_FIELDS = `
ctx.partition_name,
ctx.screenshot_path,
ctx.artifacts_json,
ctx.state_json
ctx.state_json,
kcts.ticket_status AS kuaishou_ticket_status,
kcts.ticket_code_masked AS kuaishou_ticket_code_masked,
kcts.bind_url AS kuaishou_bind_url,
kcts.vn_phone AS kuaishou_vn_phone,
kcts.role_name AS kuaishou_role_name,
kcts.role_id AS kuaishou_role_id,
kcts.dispatch_status AS kuaishou_dispatch_status,
kcts.consume_status AS kuaishou_consume_status,
kcts.source_key AS kuaishou_source_key,
kcts.updated_at AS kuaishou_state_updated_at
`
const TASK_JOINS = `
FROM fulfillment_tasks ft
LEFT JOIN claim_tokens ct ON ct.task_id = ft.id
LEFT JOIN task_runtime_contexts ctx ON ctx.task_id = ft.id
LEFT JOIN kuaishou_cloud_task_states kcts ON kcts.task_id = ft.id
`
function buildTaskSelect(extraFields = ''): string {
@@ -138,6 +150,13 @@ export async function createTask(input: TaskCreateInput): Promise<TaskRow | null
await upsertTaskRuntimeContextWithClient(client, taskId, input.runtimeContext, input.createdAt)
}
await syncKuaishouCloudTaskStateForTask({
taskId,
executorKey: input.executorKey,
contextJson: input.contextJson || '{}',
updatedAt: input.updatedAt,
}, client.query.bind(client) as TaskQueryExecutor)
return getTaskByIdWithExecutor(client.query.bind(client) as TaskQueryExecutor, taskId)
})
}
@@ -210,6 +229,13 @@ export async function updateTask(taskId: number | string, patch: TaskUpdatePatch
await upsertTaskRuntimeContextWithClient(client, Number(taskId), runtimeContextPatch, patch.updated_at || current.updated_at)
}
await syncKuaishouCloudTaskStateForTask({
taskId,
executorKey: next.executor_key,
contextJson: next.context_json,
updatedAt: next.updated_at,
}, client.query.bind(client) as TaskQueryExecutor)
return getTaskByIdWithExecutor(client.query.bind(client) as TaskQueryExecutor, taskId)
})
}
@@ -220,47 +246,62 @@ export async function updateTaskStatusIfCurrent(
patch: TaskUpdatePatch,
): Promise<TaskRow | null> {
const now = String(patch.updated_at || '').trim()
const result = await query<{ id: number }>(
`
UPDATE fulfillment_tasks
SET
task_status = $1,
delivery_status = COALESCE($2, delivery_status),
result_code = COALESCE($3, result_code),
result_message = COALESCE($4, result_message),
user_action_status = COALESCE($5, user_action_status),
last_error = COALESCE($6, last_error),
context_json = COALESCE($7::jsonb, context_json),
redeemed_at = COALESCE($8, redeemed_at),
updated_at = COALESCE($9, updated_at)
WHERE id = $10
AND task_status = $11
RETURNING id
`,
[
patch.task_status,
patch.delivery_status ?? null,
patch.result_code ?? null,
patch.result_message ?? null,
patch.user_action_status ?? null,
patch.last_error ?? null,
typeof patch.context_json === 'undefined'
? null
: typeof patch.context_json === 'string'
? patch.context_json
: JSON.stringify(patch.context_json || {}),
patch.redeemed_at ?? null,
now || null,
Number(taskId),
currentStatus,
],
)
return withTransaction(async (client) => {
const result = await client.query<{
id: number
executor_key: string
context_json: string | Record<string, unknown>
updated_at: string
}>(
`
UPDATE fulfillment_tasks
SET
task_status = $1,
delivery_status = COALESCE($2, delivery_status),
result_code = COALESCE($3, result_code),
result_message = COALESCE($4, result_message),
user_action_status = COALESCE($5, user_action_status),
last_error = COALESCE($6, last_error),
context_json = COALESCE($7::jsonb, context_json),
redeemed_at = COALESCE($8, redeemed_at),
updated_at = COALESCE($9, updated_at)
WHERE id = $10
AND task_status = $11
RETURNING id, executor_key, context_json, updated_at
`,
[
patch.task_status,
patch.delivery_status ?? null,
patch.result_code ?? null,
patch.result_message ?? null,
patch.user_action_status ?? null,
patch.last_error ?? null,
typeof patch.context_json === 'undefined'
? null
: typeof patch.context_json === 'string'
? patch.context_json
: JSON.stringify(patch.context_json || {}),
patch.redeemed_at ?? null,
now || null,
Number(taskId),
currentStatus,
],
)
if (!result.rows[0]) {
return null
}
const updated = result.rows[0]
if (!updated) {
return null
}
return getTaskById(taskId)
await syncKuaishouCloudTaskStateForTask({
taskId: updated.id,
executorKey: updated.executor_key,
contextJson: updated.context_json,
updatedAt: updated.updated_at,
}, client.query.bind(client) as TaskQueryExecutor)
return getTaskByIdWithExecutor(client.query.bind(client) as TaskQueryExecutor, taskId)
})
}
export async function getTaskById(taskId: number | string): Promise<TaskRow | null> {
@@ -315,7 +356,7 @@ export async function listTasks({
if (roleId) {
params.push(`%${roleId}%`)
filters.push(`ctx.role_id ILIKE $${params.length}`)
filters.push(`COALESCE(NULLIF(kcts.role_id, ''), ctx.role_id) ILIKE $${params.length}`)
}
if (dateFrom) {
@@ -335,6 +376,7 @@ export async function listTasks({
FROM fulfillment_tasks ft
LEFT JOIN order_items oi ON oi.id = ft.order_item_id
LEFT JOIN task_runtime_contexts ctx ON ctx.task_id = ft.id
LEFT JOIN kuaishou_cloud_task_states kcts ON kcts.task_id = ft.id
${whereClause}
`,
params,
@@ -1,10 +1,26 @@
import { query } from '../../db/client.js'
import { TASK_STATUS } from '../../domain/task-status.js'
import { nowIso } from '../../utils/time.js'
type JsonObject = Record<string, any>
export async function getAdminDashboardSummary() {
const todayPrefix = nowIso().slice(0, 10)
const paidPendingClaimStatuses = [
TASK_STATUS.PAID,
TASK_STATUS.LINK_GENERATED,
TASK_STATUS.PENDING_BINDING_PREPARE,
TASK_STATUS.WAITING_BINDING,
]
const claimingStatuses = [
TASK_STATUS.CLAIMED,
TASK_STATUS.ROLE_CONFIRMED,
TASK_STATUS.REDEEMING,
]
const abnormalStatuses = [
TASK_STATUS.RETRY_PENDING,
TASK_STATUS.MANUAL_REVIEW,
]
const result = await query(
`
SELECT
@@ -12,26 +28,32 @@ export async function getAdminDashboardSummary() {
(
SELECT COUNT(*)::int
FROM fulfillment_tasks
WHERE task_status IN ('paid', 'link_generated', 'pending_binding_prepare', 'waiting_binding')
WHERE task_status = ANY($2::text[])
) AS paid_pending_claim,
(
SELECT COUNT(*)::int
FROM fulfillment_tasks
WHERE task_status IN ('claimed', 'role_confirmed', 'redeeming')
WHERE task_status = ANY($3::text[])
) AS claiming_tasks,
(
SELECT COUNT(*)::int
FROM fulfillment_tasks
WHERE task_status = 'redeemed'
WHERE task_status = $4
AND to_char(updated_at AT TIME ZONE 'Asia/Shanghai', 'YYYY-MM-DD') = $1
) AS redeemed_today,
(
SELECT COUNT(*)::int
FROM fulfillment_tasks
WHERE task_status IN ('retry_pending', 'manual_review')
WHERE task_status = ANY($5::text[])
) AS abnormal_tasks
`,
[todayPrefix],
[
todayPrefix,
paidPendingClaimStatuses,
claimingStatuses,
TASK_STATUS.REDEEMED,
abnormalStatuses,
],
)
const summary: JsonObject = result.rows[0] || {}
@@ -23,6 +23,7 @@ import {
getTaskPrimaryClaimTokenId,
isManualDispatchTask,
mapKuaishouCloudFulfillmentContext,
mapKuaishouCloudTaskStateProjection,
mapManualDispatchContext,
mapRedeemResolutionContext,
parseTaskContext,
@@ -180,10 +181,11 @@ export async function getAdminTaskDetail(
const viewerContext = createAdminViewerContext(session)
const claimUrl = claimToken ? buildClaimUrl(claimToken.token) : ''
const screenshotUrl = await resolveAdminTaskScreenshotUrl(task, viewerContext)
const cloudSourceLabelMap = buildCloudSourceLabelMap()
const kuaishouCloudFulfillment = mapKuaishouCloudFulfillmentContext(
taskContext.kuaishouCloudFulfillment,
{ cloudSourceLabelMap: buildCloudSourceLabelMap() },
)
{ cloudSourceLabelMap },
) || mapKuaishouCloudTaskStateProjection(task, { cloudSourceLabelMap })
return {
task: mapAdminTaskListItem({
@@ -160,6 +160,105 @@ export function mapKuaishouCloudFulfillmentContext(
}
}
export function mapKuaishouCloudTaskStateProjection(
task: TaskLike | null | undefined,
options: KuaishouCloudFulfillmentMapOptions = {},
): JsonRecord | null {
if (!task || String(task.executor_key || '').trim() !== 'kuaishou_ct_assisted') {
return null
}
const hasProjection = [
task.kuaishou_ticket_status,
task.kuaishou_ticket_code_masked,
task.kuaishou_bind_url,
task.kuaishou_vn_phone,
task.kuaishou_role_name,
task.kuaishou_role_id,
task.kuaishou_dispatch_status,
task.kuaishou_consume_status,
task.kuaishou_source_key,
].some((value) => String(value || '').trim())
if (!hasProjection) {
return null
}
const sourceKey = String(task.kuaishou_source_key || '').trim()
const roleName = String(task.kuaishou_role_name || '').trim()
const roleId = String(task.kuaishou_role_id || '').trim()
return {
flowType: 'kuaishou_cloud_fulfillment',
configId: '',
internalSkuCode: '',
internalSkuName: '',
ticket: {
code: String(task.kuaishou_ticket_code_masked || '').trim(),
capturedAt: null,
status: String(task.kuaishou_ticket_status || 'pending').trim() || 'pending',
verifiedAt: null,
goodsTitle: '',
leftCount: 0,
},
binding: {
prepareStatus: task.kuaishou_bind_url ? 'ready' : 'pending',
cloudSourceKeys: sourceKey ? [sourceKey] : [],
cloudSourceLabels: sourceKey ? [resolveCloudSourceLabel(sourceKey, options.cloudSourceLabelMap)] : [],
resolvedSourceKey: sourceKey,
resolvedSourceLabel: sourceKey ? resolveCloudSourceLabel(sourceKey, options.cloudSourceLabelMap) : '',
skuId: 0,
skuName: '',
vnKey: '',
vnId: 0,
vnPhone: String(task.kuaishou_vn_phone || '').trim(),
bindUrl: String(task.kuaishou_bind_url || '').trim(),
bindPreparedAt: null,
bindExpiresAt: null,
bindProbeAt: null,
bindProbeStatus: '',
bindProbeMessage: '',
},
role: {
status: roleName || roleId ? 'ready' : 'pending',
name: roleName,
rid: roleId,
refreshedAt: null,
errorMessage: '',
rawInfo: null,
},
purchase: {
autoBuyEnabled: true,
minAssetReserve: 0,
usedKnapsack: false,
purchaseTriggered: false,
assetBefore: 0,
assetAfter: 0,
purchaseAt: null,
},
dispatch: {
status: String(task.kuaishou_dispatch_status || 'pending').trim() || 'pending',
dispatchAt: null,
sendType: 0,
note: '',
},
returnNumber: {
status: 'pending',
returnedAt: null,
autoReturnEnabled: false,
},
consume: {
status: String(task.kuaishou_consume_status || 'pending').trim() || 'pending',
shopId: '',
shopName: '',
autoConsumeEnabled: false,
consumedAt: null,
errorMessage: '',
},
notes: '',
}
}
function resolveCloudSourceLabel(sourceKey: string, labelMap: CloudSourceLabelMap | undefined) {
const key = String(sourceKey || '').trim()
if (!key) return ''
@@ -102,8 +102,8 @@ export function mapAdminTaskListItem(
resourceStatus: fulfillment.resourceStatus,
customerStatus: fulfillment.customerStatus,
loginType: task.login_type,
roleName: task.role_name,
roleId: task.role_id,
roleName: resolveTaskDisplayRoleName(task),
roleId: resolveTaskDisplayRoleId(task),
runtimeSessionId: task.runtime_session_id,
claimedAt: task.claimed_at,
roleConfirmedAt: task.role_confirmed_at,
@@ -269,3 +269,11 @@ function normalizeRecord(value: unknown): JsonRecord {
function getTaskRetryCount(task: TaskRow | null | undefined): number {
return Number(task?.attempt_count || 0)
}
function resolveTaskDisplayRoleName(task: TaskRow): string {
return String(task.kuaishou_role_name || task.role_name || '').trim()
}
function resolveTaskDisplayRoleId(task: TaskRow): string {
return String(task.kuaishou_role_id || task.role_id || '').trim()
}
@@ -20,6 +20,7 @@ import {
} from '../../platforms/kuaishou-eticket/source-config-service.js'
import { createHttpError } from '../../../utils/http.js'
import { nowIso } from '../../../utils/time.js'
import { TASK_STATUS } from '../../../domain/task-status.js'
import {
createAdminViewerContext,
parseTaskContext,
@@ -200,7 +201,7 @@ export async function prepareAdminTaskKuaishouCloudFulfillment(
}
const updatedTask = await updateTask(task.id, {
task_status: 'waiting_binding',
task_status: TASK_STATUS.WAITING_BINDING,
user_action_status: 'pending',
claim_token: claimLinkState.token || task.claim_token || '',
claim_expires_at: claimLinkState.expiredAt || getTaskClaimExpiresAt(task),
@@ -334,7 +335,7 @@ export async function dispatchAdminTaskKuaishouCloudFulfillment(
}
const updatedTask = await updateTask(task.id, {
task_status: 'dispatched_pending_return',
task_status: TASK_STATUS.DISPATCHED_PENDING_RETURN,
delivery_status: 'delivered',
result_code: 'kuaishou_cloud_dispatched',
result_message: String(dispatchResult.responseMessage || dispatchResult.note || 'cloudtentacles 发货成功').trim(),
@@ -3,6 +3,7 @@ import { updateTask } from '../../../repositories/task-repo.js'
import { createTaskEvent } from '../../../repositories/task-event-repo.js'
import { createHttpError } from '../../../utils/http.js'
import { nowIso } from '../../../utils/time.js'
import { TASK_STATUS } from '../../../domain/task-status.js'
import {
canViewerCloseTask,
createAdminViewerContext,
@@ -44,7 +45,7 @@ export async function closeAdminTask(
}
const updatedTask = await updateTask(task.id, {
task_status: 'closed',
task_status: TASK_STATUS.CLOSED,
delivery_status: 'closed',
result_code: 'closed_by_admin',
result_message: '任务已关闭,领取链接已失效',
+10
View File
@@ -70,6 +70,16 @@ export type TaskRow = {
retry_count: number
created_at: string
updated_at: string
kuaishou_ticket_status?: string
kuaishou_ticket_code_masked?: string
kuaishou_bind_url?: string
kuaishou_vn_phone?: string
kuaishou_role_name?: string
kuaishou_role_id?: string
kuaishou_dispatch_status?: string
kuaishou_consume_status?: string
kuaishou_source_key?: string
kuaishou_state_updated_at?: string | null
claimed_at: string | null
role_confirmed_at: string | null
redeemed_at: string | null