From 5a195e21502c37f0242619d4629b51c589e52c95 Mon Sep 17 00:00:00 2001 From: yml2213 Date: Fri, 21 Aug 2026 18:11:21 +0800 Subject: [PATCH] =?UTF-8?q?=E6=8B=86=E5=88=86=E6=89=93=E6=89=8B=E6=8E=92?= =?UTF-8?q?=E5=90=8D=E4=BB=93=E5=82=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../src/repositories/worker-platform/index.ts | 1 + .../worker-platform/work-order-repo.ts | 2 +- .../worker-platform/worker-ranking-repo.ts | 279 ++++++++++++++++++ .../worker-platform/worker-repo.ts | 267 +---------------- 4 files changed, 287 insertions(+), 262 deletions(-) create mode 100644 apps/backend/src/repositories/worker-platform/worker-ranking-repo.ts diff --git a/apps/backend/src/repositories/worker-platform/index.ts b/apps/backend/src/repositories/worker-platform/index.ts index 94efc081..0630a972 100644 --- a/apps/backend/src/repositories/worker-platform/index.ts +++ b/apps/backend/src/repositories/worker-platform/index.ts @@ -5,6 +5,7 @@ export * from './worker-session-repo.js' export * from './worker-level-repo.js' export * from './worker-wallet-repo.js' export * from './worker-finance-repo.js' +export * from './worker-ranking-repo.js' export * from './work-order-repo.js' export * from './work-order-event-repo.js' export * from './work-order-share-repo.js' diff --git a/apps/backend/src/repositories/worker-platform/work-order-repo.ts b/apps/backend/src/repositories/worker-platform/work-order-repo.ts index 3bb4fdb0..0b8264ea 100644 --- a/apps/backend/src/repositories/worker-platform/work-order-repo.ts +++ b/apps/backend/src/repositories/worker-platform/work-order-repo.ts @@ -20,7 +20,7 @@ import type { WorkOrderStatistics, } from './types.js' import type { PoolClient } from 'pg' -import { maybeUpgradeWorkerLevelWithClient } from './worker-repo.js' +import { maybeUpgradeWorkerLevelWithClient } from './worker-ranking-repo.js' export const WORK_ORDER_SELECT = ` SELECT diff --git a/apps/backend/src/repositories/worker-platform/worker-ranking-repo.ts b/apps/backend/src/repositories/worker-platform/worker-ranking-repo.ts new file mode 100644 index 00000000..1418f963 --- /dev/null +++ b/apps/backend/src/repositories/worker-platform/worker-ranking-repo.ts @@ -0,0 +1,279 @@ +import { query } from '../../db/client.js' +import type { WorkerUserRow } from './types.js' +import type { PoolClient } from 'pg' + +const WORKER_USER_SELECT = ` + SELECT + wu.*, + wl.level_key, + wl.name AS level_name, + wl.permission_json AS level_permission_json, + inv.username AS inviter_username, + inv.display_name AS inviter_display_name, + ww.available_amount, + ww.frozen_deposit_amount, + ww.pending_unfreeze_amount, + ww.total_credited_amount, + ww.total_settled_amount + FROM worker_users wu + LEFT JOIN worker_levels wl ON wl.id = wu.level_id + LEFT JOIN worker_users inv ON inv.id = wu.inviter_id + LEFT JOIN worker_wallets ww ON ww.worker_id = wu.id +` + +/** 完成单数按源订单去重,拼单份数只影响收益,不增加完成单数。 */ +export async function countWorkerAcceptedOrders(workerId: number | string): Promise { + const result = await query<{ total: number }>(buildWorkerAcceptedOrdersCountSql(), [ + Number(workerId), + ]) + return Number(result.rows[0]?.total || 0) +} + +export async function listWorkerCompletionLeaderboard(input: { + limit: number + currentWorkerId: number + period?: 'all' | 'daily' +}): Promise< + Array<{ + rank: number + workerId: number + displayName: string + levelName: string + acceptedOrderCount: number + }> +> { + const limit = Math.min(50, Math.max(1, Math.floor(Number(input.limit || 10)))) + const dailyOnly = input.period === 'daily' + const result = await query<{ + rank: number + worker_id: number + display_name: string + level_name: string | null + accepted_order_count: number + }>( + ` + WITH leaderboard_day AS ( + SELECT + date_trunc('day', CURRENT_TIMESTAMP AT TIME ZONE 'Asia/Shanghai') + AT TIME ZONE 'Asia/Shanghai' AS starts_at + ), + accepted_totals AS ( + SELECT + worker_id, + COUNT(DISTINCT order_key)::int AS accepted_order_count, + MAX(last_accepted_at) AS last_accepted_at + FROM ( + SELECT + assigned_worker_id AS worker_id, + COALESCE( + 'order:' || NULLIF(order_id::text, ''), + 'platform_order:' || NULLIF(BTRIM(platform_order_id), ''), + 'work_order:' || id::text + ) AS order_key, + MAX(accepted_at) AS last_accepted_at + FROM work_orders + CROSS JOIN leaderboard_day + WHERE assigned_worker_id IS NOT NULL + AND status = 'accepted' + AND ( + $3::boolean = false + OR (accepted_at >= leaderboard_day.starts_at AND accepted_at < leaderboard_day.starts_at + INTERVAL '1 day') + ) + GROUP BY assigned_worker_id, order_key + + UNION ALL + + SELECT + wos.worker_id, + COALESCE( + 'order:' || NULLIF(wo.order_id::text, ''), + 'platform_order:' || NULLIF(BTRIM(wo.platform_order_id), ''), + 'work_order:' || wo.id::text + ) AS order_key, + MAX(wos.accepted_at) AS last_accepted_at + FROM work_order_shares wos + INNER JOIN work_orders wo ON wo.id = wos.work_order_id + CROSS JOIN leaderboard_day + WHERE wos.status = 'accepted' + AND ( + $3::boolean = false + OR (wos.accepted_at >= leaderboard_day.starts_at AND wos.accepted_at < leaderboard_day.starts_at + INTERVAL '1 day') + ) + GROUP BY wos.worker_id, order_key + ) AS accepted_records + GROUP BY worker_id + ), + ranked_workers AS ( + SELECT + ROW_NUMBER() OVER ( + ORDER BY totals.accepted_order_count DESC, totals.last_accepted_at DESC NULLS LAST, wu.id ASC + )::int AS rank, + wu.id AS worker_id, + COALESCE(NULLIF(wu.display_name, ''), wu.username) AS display_name, + wl.name AS level_name, + totals.accepted_order_count + FROM worker_users wu + INNER JOIN accepted_totals totals ON totals.worker_id = wu.id + LEFT JOIN worker_levels wl ON wl.id = wu.level_id + WHERE wu.status = 'active' + ) + SELECT rank, worker_id, display_name, level_name, accepted_order_count + FROM ranked_workers + WHERE rank <= $1 OR worker_id = $2 + ORDER BY rank ASC + `, + [limit, Number(input.currentWorkerId), dailyOnly], + ) + return result.rows.map((row) => ({ + rank: Number(row.rank || 0), + workerId: Number(row.worker_id || 0), + displayName: String(row.display_name || '').trim(), + levelName: String(row.level_name || '').trim(), + acceptedOrderCount: Number(row.accepted_order_count || 0), + })) +} + +export async function maybeUpgradeWorkerLevelWithClient( + client: PoolClient, + workerId: number, + now: string, +): Promise { + const current = await client.query<{ + id: number + level_id: number | null + level_penalty_accepted_count: number + upgrade_threshold: number | null + }>( + ` + SELECT + wu.id, + wu.level_id, + wu.level_penalty_accepted_count, + COALESCE((wl.permission_json->>'upgradeThreshold')::int, 0) AS upgrade_threshold + FROM worker_users wu + LEFT JOIN worker_levels wl ON wl.id = wu.level_id + WHERE wu.id = $1 + LIMIT 1 + FOR UPDATE OF wu + `, + [workerId], + ) + const currentRow = current.rows[0] + if (!currentRow) return null + + const acceptedCount = await countWorkerAcceptedOrdersWithClient(client, workerId) + const effectiveAcceptedCount = Math.max( + 0, + acceptedCount - Number(currentRow.level_penalty_accepted_count || 0), + ) + const nextLevelResult = await client.query<{ + id: number + upgrade_threshold: number + }>( + ` + SELECT + id, + COALESCE((permission_json->>'upgradeThreshold')::int, 0) AS upgrade_threshold + FROM worker_levels + WHERE status = 'active' + AND COALESCE((permission_json->>'upgradeThreshold')::int, 0) <= $1 + ORDER BY COALESCE((permission_json->>'upgradeThreshold')::int, 0) DESC, sort_order ASC + LIMIT 1 + `, + [effectiveAcceptedCount], + ) + const nextLevel = nextLevelResult.rows[0] + const nextLevelId = Number(nextLevel?.id || 0) + if (!nextLevelId) return null + if ( + !shouldApplyAutomaticWorkerLevel({ + currentLevelId: currentRow.level_id, + currentThreshold: Number(currentRow.upgrade_threshold || 0), + candidateLevelId: nextLevelId, + candidateThreshold: Number(nextLevel?.upgrade_threshold || 0), + }) + ) { + return null + } + + await client.query('UPDATE worker_users SET level_id = $1, updated_at = $2 WHERE id = $3', [ + nextLevelId, + now, + workerId, + ]) + await client.query( + ` + INSERT INTO worker_level_logs ( + worker_id, from_level_id, to_level_id, reason, payload_json, created_at + ) VALUES ($1, $2, $3, 'auto', $4::jsonb, $5) + `, + [ + workerId, + currentRow.level_id || null, + nextLevelId, + JSON.stringify({ acceptedCount, effectiveAcceptedCount }), + now, + ], + ) + return getWorkerUserByIdWithClient(client, workerId) +} + +export function shouldApplyAutomaticWorkerLevel(input: { + currentLevelId: number | null + currentThreshold: number + candidateLevelId: number + candidateThreshold: number +}) { + if (!input.currentLevelId) return true + if (Number(input.currentLevelId) === Number(input.candidateLevelId)) return false + return Number(input.candidateThreshold) > Number(input.currentThreshold) +} + +async function countWorkerAcceptedOrdersWithClient( + client: PoolClient, + workerId: number, +): Promise { + const result = await client.query<{ total: number }>(buildWorkerAcceptedOrdersCountSql(), [ + workerId, + ]) + return Number(result.rows[0]?.total || 0) +} + +async function getWorkerUserByIdWithClient( + client: PoolClient, + workerId: number, +): Promise { + const result = await client.query( + `${WORKER_USER_SELECT} WHERE wu.id = $1 LIMIT 1`, + [workerId], + ) + return result.rows[0] || null +} + +function buildWorkerAcceptedOrdersCountSql() { + return ` + SELECT COUNT(DISTINCT order_key)::int AS total + FROM ( + SELECT + COALESCE( + 'order:' || NULLIF(wo.order_id::text, ''), + 'platform_order:' || NULLIF(BTRIM(wo.platform_order_id), ''), + 'work_order:' || wo.id::text + ) AS order_key + FROM work_orders wo + WHERE wo.assigned_worker_id = $1 AND wo.status = 'accepted' + + UNION + + SELECT + COALESCE( + 'order:' || NULLIF(wo.order_id::text, ''), + 'platform_order:' || NULLIF(BTRIM(wo.platform_order_id), ''), + 'work_order:' || wo.id::text + ) AS order_key + FROM work_order_shares wos + INNER JOIN work_orders wo ON wo.id = wos.work_order_id + WHERE wos.worker_id = $1 AND wos.status = 'accepted' + ) AS completed_orders + ` +} diff --git a/apps/backend/src/repositories/worker-platform/worker-repo.ts b/apps/backend/src/repositories/worker-platform/worker-repo.ts index f4fe8b74..1a261af8 100644 --- a/apps/backend/src/repositories/worker-platform/worker-repo.ts +++ b/apps/backend/src/repositories/worker-platform/worker-repo.ts @@ -275,267 +275,12 @@ export { upsertWorkerWithdrawalAccount, } from './worker-finance-repo.js' -/** 完成单数按源订单去重,拼单份数只影响收益,不增加完成单数。 */ -export async function countWorkerAcceptedOrders(workerId: number | string): Promise { - const result = await query<{ total: number }>( - ` - SELECT COUNT(DISTINCT order_key)::int AS total - FROM ( - SELECT - COALESCE( - 'order:' || NULLIF(wo.order_id::text, ''), - 'platform_order:' || NULLIF(BTRIM(wo.platform_order_id), ''), - 'work_order:' || wo.id::text - ) AS order_key - FROM work_orders wo - WHERE wo.assigned_worker_id = $1 AND wo.status = 'accepted' - - UNION - - SELECT - COALESCE( - 'order:' || NULLIF(wo.order_id::text, ''), - 'platform_order:' || NULLIF(BTRIM(wo.platform_order_id), ''), - 'work_order:' || wo.id::text - ) AS order_key - FROM work_order_shares wos - INNER JOIN work_orders wo ON wo.id = wos.work_order_id - WHERE wos.worker_id = $1 AND wos.status = 'accepted' - ) AS completed_orders - `, - [Number(workerId)], - ) - return Number(result.rows[0]?.total || 0) -} - -export async function listWorkerCompletionLeaderboard(input: { - limit: number - currentWorkerId: number - period?: 'all' | 'daily' -}): Promise< - Array<{ - rank: number - workerId: number - displayName: string - levelName: string - acceptedOrderCount: number - }> -> { - const limit = Math.min(50, Math.max(1, Math.floor(Number(input.limit || 10)))) - const dailyOnly = input.period === 'daily' - const result = await query<{ - rank: number - worker_id: number - display_name: string - level_name: string | null - accepted_order_count: number - }>( - ` - WITH leaderboard_day AS ( - SELECT - date_trunc('day', CURRENT_TIMESTAMP AT TIME ZONE 'Asia/Shanghai') - AT TIME ZONE 'Asia/Shanghai' AS starts_at - ), - accepted_totals AS ( - SELECT - worker_id, - COUNT(DISTINCT order_key)::int AS accepted_order_count, - MAX(last_accepted_at) AS last_accepted_at - FROM ( - SELECT - assigned_worker_id AS worker_id, - COALESCE( - 'order:' || NULLIF(order_id::text, ''), - 'platform_order:' || NULLIF(BTRIM(platform_order_id), ''), - 'work_order:' || id::text - ) AS order_key, - MAX(accepted_at) AS last_accepted_at - FROM work_orders - CROSS JOIN leaderboard_day - WHERE assigned_worker_id IS NOT NULL - AND status = 'accepted' - AND ( - $3::boolean = false - OR (accepted_at >= leaderboard_day.starts_at AND accepted_at < leaderboard_day.starts_at + INTERVAL '1 day') - ) - GROUP BY assigned_worker_id, order_key - - UNION ALL - - SELECT - wos.worker_id, - COALESCE( - 'order:' || NULLIF(wo.order_id::text, ''), - 'platform_order:' || NULLIF(BTRIM(wo.platform_order_id), ''), - 'work_order:' || wo.id::text - ) AS order_key, - MAX(wos.accepted_at) AS last_accepted_at - FROM work_order_shares wos - INNER JOIN work_orders wo ON wo.id = wos.work_order_id - CROSS JOIN leaderboard_day - WHERE wos.status = 'accepted' - AND ( - $3::boolean = false - OR (wos.accepted_at >= leaderboard_day.starts_at AND wos.accepted_at < leaderboard_day.starts_at + INTERVAL '1 day') - ) - GROUP BY wos.worker_id, order_key - ) AS accepted_records - GROUP BY worker_id - ), - ranked_workers AS ( - SELECT - ROW_NUMBER() OVER ( - ORDER BY totals.accepted_order_count DESC, totals.last_accepted_at DESC NULLS LAST, wu.id ASC - )::int AS rank, - wu.id AS worker_id, - COALESCE(NULLIF(wu.display_name, ''), wu.username) AS display_name, - wl.name AS level_name, - totals.accepted_order_count - FROM worker_users wu - INNER JOIN accepted_totals totals ON totals.worker_id = wu.id - LEFT JOIN worker_levels wl ON wl.id = wu.level_id - WHERE wu.status = 'active' - ) - SELECT rank, worker_id, display_name, level_name, accepted_order_count - FROM ranked_workers - WHERE rank <= $1 OR worker_id = $2 - ORDER BY rank ASC - `, - [limit, Number(input.currentWorkerId), dailyOnly], - ) - return result.rows.map((row) => ({ - rank: Number(row.rank || 0), - workerId: Number(row.worker_id || 0), - displayName: String(row.display_name || '').trim(), - levelName: String(row.level_name || '').trim(), - acceptedOrderCount: Number(row.accepted_order_count || 0), - })) -} - -export async function maybeUpgradeWorkerLevelWithClient( - client: PoolClient, - workerId: number, - now: string, -): Promise { - const current = await client.query<{ - id: number - level_id: number | null - level_penalty_accepted_count: number - upgrade_threshold: number | null - }>( - ` - SELECT - wu.id, - wu.level_id, - wu.level_penalty_accepted_count, - COALESCE((wl.permission_json->>'upgradeThreshold')::int, 0) AS upgrade_threshold - FROM worker_users wu - LEFT JOIN worker_levels wl ON wl.id = wu.level_id - WHERE wu.id = $1 - LIMIT 1 - FOR UPDATE OF wu - `, - [workerId], - ) - const currentRow = current.rows[0] - if (!currentRow) return null - - const acceptedResult = await client.query<{ total: number }>( - ` - SELECT COUNT(DISTINCT order_key)::int AS total - FROM ( - SELECT - COALESCE( - 'order:' || NULLIF(wo.order_id::text, ''), - 'platform_order:' || NULLIF(BTRIM(wo.platform_order_id), ''), - 'work_order:' || wo.id::text - ) AS order_key - FROM work_orders wo - WHERE wo.assigned_worker_id = $1 AND wo.status = 'accepted' - - UNION - - SELECT - COALESCE( - 'order:' || NULLIF(wo.order_id::text, ''), - 'platform_order:' || NULLIF(BTRIM(wo.platform_order_id), ''), - 'work_order:' || wo.id::text - ) AS order_key - FROM work_order_shares wos - INNER JOIN work_orders wo ON wo.id = wos.work_order_id - WHERE wos.worker_id = $1 AND wos.status = 'accepted' - ) AS completed_orders - `, - [workerId], - ) - const acceptedCount = Number(acceptedResult.rows[0]?.total || 0) - const effectiveAcceptedCount = Math.max( - 0, - acceptedCount - Number(currentRow.level_penalty_accepted_count || 0), - ) - - const nextLevelResult = await client.query<{ - id: number - upgrade_threshold: number - }>( - ` - SELECT - id, - COALESCE((permission_json->>'upgradeThreshold')::int, 0) AS upgrade_threshold - FROM worker_levels - WHERE status = 'active' - AND COALESCE((permission_json->>'upgradeThreshold')::int, 0) <= $1 - ORDER BY COALESCE((permission_json->>'upgradeThreshold')::int, 0) DESC, sort_order ASC - LIMIT 1 - `, - [effectiveAcceptedCount], - ) - const nextLevel = nextLevelResult.rows[0] - const nextLevelId = Number(nextLevel?.id || 0) - if (!nextLevelId) return null - if ( - !shouldApplyAutomaticWorkerLevel({ - currentLevelId: currentRow.level_id, - currentThreshold: Number(currentRow.upgrade_threshold || 0), - candidateLevelId: nextLevelId, - candidateThreshold: Number(nextLevel?.upgrade_threshold || 0), - }) - ) { - return null - } - - await client.query('UPDATE worker_users SET level_id = $1, updated_at = $2 WHERE id = $3', [ - nextLevelId, - now, - workerId, - ]) - await client.query( - ` - INSERT INTO worker_level_logs ( - worker_id, from_level_id, to_level_id, reason, payload_json, created_at - ) VALUES ($1, $2, $3, 'auto', $4::jsonb, $5) - `, - [ - workerId, - currentRow.level_id || null, - nextLevelId, - JSON.stringify({ acceptedCount, effectiveAcceptedCount }), - now, - ], - ) - return getWorkerUserByIdWithClient(client, workerId) -} - -export function shouldApplyAutomaticWorkerLevel(input: { - currentLevelId: number | null - currentThreshold: number - candidateLevelId: number - candidateThreshold: number -}) { - if (!input.currentLevelId) return true - if (Number(input.currentLevelId) === Number(input.candidateLevelId)) return false - return Number(input.candidateThreshold) > Number(input.currentThreshold) -} +export { + countWorkerAcceptedOrders, + listWorkerCompletionLeaderboard, + maybeUpgradeWorkerLevelWithClient, + shouldApplyAutomaticWorkerLevel, +} from './worker-ranking-repo.js' async function getWorkerUserByIdWithClient( client: PoolClient,