import { query } from '../db/client.js' import { TASK_STATUS } from '../domain/task-status.js' import { maskCode } from '../utils/masking.js' import { parseTaskContext } from '../utils/task-json.js' import { normalizeTimestampIso } from '../utils/time.js' import type { JsonObject, JsonRecord } from '../types/json.js' type QueryExecutor = (text: string, params?: unknown[]) => Promise<{ rows: unknown[] }> export type KuaishouCloudSourceLoadStat = { sourceKey: string activeCount: number redeemingCount: number } export type KuaishouCloudTaskStateSyncInput = { taskId: number | string executorKey?: unknown contextJson?: unknown updatedAt?: unknown } export async function syncKuaishouCloudTaskStateForTask( input: KuaishouCloudTaskStateSyncInput, executor: QueryExecutor = query, ): Promise { 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 = normalizeTimestampIso(input.updatedAt) 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, ], ) } export async function listKuaishouCloudSourceLoadStats( sourceKeys: unknown[] = [], executor: QueryExecutor = query, ): Promise { const candidates = normalizeStringArray(sourceKeys) if (candidates.length === 0) { return [] } const activeStatuses = [ TASK_STATUS.WAITING_BINDING, TASK_STATUS.ROLE_CONFIRMED, TASK_STATUS.REDEEMING, TASK_STATUS.DISPATCHED_PENDING_RETURN, ] const result = await executor( ` SELECT kcts.source_key, COUNT(*) FILTER (WHERE ft.task_status = ANY($2::text[])) AS active_count, COUNT(*) FILTER (WHERE ft.task_status = $3) AS redeeming_count FROM fulfillment_tasks ft JOIN kuaishou_cloud_task_states kcts ON kcts.task_id = ft.id WHERE ft.executor_key = 'kuaishou_ct_assisted' AND kcts.source_key = ANY($1::text[]) AND ft.task_status = ANY($2::text[]) GROUP BY kcts.source_key `, [candidates, activeStatuses, TASK_STATUS.REDEEMING], ) return result.rows .map((row) => { const source = asRecord(row) return { sourceKey: String(source.source_key || '').trim(), activeCount: Number(source.active_count || 0) || 0, redeemingCount: Number(source.redeeming_count || 0) || 0, } }) .filter((row) => row.sourceKey) } 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 normalizeStringArray(value: unknown): string[] { if (Array.isArray(value)) { return Array.from(new Set(value.map((item) => String(item || '').trim()).filter(Boolean))) } if (typeof value === 'string') { return Array.from( new Set( value .split(',') .map((item) => item.trim()) .filter(Boolean), ), ) } return [] } function pickFirstNonEmpty(values: unknown[]): string { for (const value of values) { const normalized = String(value || '').trim() if (normalized) { return normalized } } return '' }