拆分打手排名仓储

This commit is contained in:
yml2213
2026-08-21 18:11:21 +08:00
parent ec29469fcb
commit 5a195e2150
4 changed files with 287 additions and 262 deletions
@@ -5,6 +5,7 @@ export * from './worker-session-repo.js'
export * from './worker-level-repo.js' export * from './worker-level-repo.js'
export * from './worker-wallet-repo.js' export * from './worker-wallet-repo.js'
export * from './worker-finance-repo.js' export * from './worker-finance-repo.js'
export * from './worker-ranking-repo.js'
export * from './work-order-repo.js' export * from './work-order-repo.js'
export * from './work-order-event-repo.js' export * from './work-order-event-repo.js'
export * from './work-order-share-repo.js' export * from './work-order-share-repo.js'
@@ -20,7 +20,7 @@ import type {
WorkOrderStatistics, WorkOrderStatistics,
} from './types.js' } from './types.js'
import type { PoolClient } from 'pg' import type { PoolClient } from 'pg'
import { maybeUpgradeWorkerLevelWithClient } from './worker-repo.js' import { maybeUpgradeWorkerLevelWithClient } from './worker-ranking-repo.js'
export const WORK_ORDER_SELECT = ` export const WORK_ORDER_SELECT = `
SELECT SELECT
@@ -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<number> {
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<WorkerUserRow | null> {
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<number> {
const result = await client.query<{ total: number }>(buildWorkerAcceptedOrdersCountSql(), [
workerId,
])
return Number(result.rows[0]?.total || 0)
}
async function getWorkerUserByIdWithClient(
client: PoolClient,
workerId: number,
): Promise<WorkerUserRow | null> {
const result = await client.query<WorkerUserRow>(
`${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
`
}
@@ -275,267 +275,12 @@ export {
upsertWorkerWithdrawalAccount, upsertWorkerWithdrawalAccount,
} from './worker-finance-repo.js' } from './worker-finance-repo.js'
/** 完成单数按源订单去重,拼单份数只影响收益,不增加完成单数。 */ export {
export async function countWorkerAcceptedOrders(workerId: number | string): Promise<number> { countWorkerAcceptedOrders,
const result = await query<{ total: number }>( listWorkerCompletionLeaderboard,
` maybeUpgradeWorkerLevelWithClient,
SELECT COUNT(DISTINCT order_key)::int AS total shouldApplyAutomaticWorkerLevel,
FROM ( } from './worker-ranking-repo.js'
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<WorkerUserRow | null> {
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)
}
async function getWorkerUserByIdWithClient( async function getWorkerUserByIdWithClient(
client: PoolClient, client: PoolClient,