From a491a281ea261c981447350f73ca7394cd52ac60 Mon Sep 17 00:00:00 2001 From: yml2213 Date: Fri, 21 Aug 2026 18:13:28 +0800 Subject: [PATCH] =?UTF-8?q?=E6=8B=86=E5=88=86=E5=B7=A5=E5=8D=95=E6=9F=A5?= =?UTF-8?q?=E8=AF=A2=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-query-repo.ts | 670 +++++++++++++++++ .../worker-platform/work-order-repo.ts | 711 +----------------- 3 files changed, 698 insertions(+), 684 deletions(-) create mode 100644 apps/backend/src/repositories/worker-platform/work-order-query-repo.ts diff --git a/apps/backend/src/repositories/worker-platform/index.ts b/apps/backend/src/repositories/worker-platform/index.ts index 0630a972..7fbb8056 100644 --- a/apps/backend/src/repositories/worker-platform/index.ts +++ b/apps/backend/src/repositories/worker-platform/index.ts @@ -7,6 +7,7 @@ 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-query-repo.js' export * from './work-order-event-repo.js' export * from './work-order-share-repo.js' export * from './work-category-repo.js' diff --git a/apps/backend/src/repositories/worker-platform/work-order-query-repo.ts b/apps/backend/src/repositories/worker-platform/work-order-query-repo.ts new file mode 100644 index 00000000..3c7cd8c6 --- /dev/null +++ b/apps/backend/src/repositories/worker-platform/work-order-query-repo.ts @@ -0,0 +1,670 @@ +import { query } from '../../db/client.js' +import type { ListInput, WorkOrderRow, WorkOrderStatistics } from './types.js' +import type { PoolClient } from 'pg' +export const WORK_ORDER_SELECT = ` + SELECT + wo.*, + COALESCE(share_summary.active_share_quantity, 0)::int AS active_share_quantity, + COALESCE(share_summary.pending_share_count, 0)::int AS pending_share_count, + wc.name AS category_name, + wu.username AS worker_username, + wu.display_name AS worker_display_name + FROM work_orders wo + LEFT JOIN work_categories wc ON wc.id = wo.category_id + LEFT JOIN worker_users wu ON wu.id = wo.assigned_worker_id + LEFT JOIN LATERAL ( + SELECT + SUM(wos.quantity) FILTER (WHERE wos.status != 'cancelled') AS active_share_quantity, + COUNT(*) FILTER (WHERE wos.status = 'joined') AS pending_share_count + FROM work_order_shares wos + WHERE wos.work_order_id = wo.id + ) share_summary ON TRUE +` + +type WorkerOrderOverviewRow = { + status: string + order_count: number + settlement_amount: number + vip_evidence_pending_count: number + vip_evidence_pending_settlement_amount: number +} + +type WorkOrderStatisticsRow = { + status: string + total_count: number + user_payment_amount: number | string + reward_amount: number | string + refund_amount: number | string +} + +/** 后台按母工单拼单参与情况计算业务状态,open 只代表大厅生命周期。 */ +function adminOperationalStatusExpression(alias = 'wo', shareSummaryAlias = '') { + if (shareSummaryAlias) { + return `CASE + WHEN ${alias}.status = 'open' + AND ${alias}.sharing_enabled IS TRUE + AND COALESCE(${shareSummaryAlias}.active_share_quantity, 0) > 0 + THEN CASE + WHEN COALESCE(${shareSummaryAlias}.active_share_quantity, 0) >= ${alias}.sharing_total_quantity + AND COALESCE(${shareSummaryAlias}.pending_share_count, 0) = 0 + THEN 'pending_acceptance' + ELSE 'in_progress' + END + ELSE ${alias}.status + END` + } + + return `CASE + WHEN ${alias}.status = 'open' + AND ${alias}.sharing_enabled IS TRUE + AND COALESCE((SELECT SUM(wos.quantity) FROM work_order_shares wos WHERE wos.work_order_id = ${alias}.id AND wos.status != 'cancelled'), 0) > 0 + THEN CASE + WHEN COALESCE((SELECT SUM(wos.quantity) FROM work_order_shares wos WHERE wos.work_order_id = ${alias}.id AND wos.status != 'cancelled'), 0) >= ${alias}.sharing_total_quantity + AND COALESCE((SELECT COUNT(*) FROM work_order_shares wos WHERE wos.work_order_id = ${alias}.id AND wos.status = 'joined'), 0) = 0 + THEN 'pending_acceptance' + ELSE 'in_progress' + END + ELSE ${alias}.status + END` +} + +/** 工单统一搜索字段,兼容当前资料结构和历史资料字段。 */ +function workOrderSearchExpression(alias = 'wo') { + return `LOWER( + COALESCE(${alias}.work_order_no, '') || ' ' || + COALESCE(${alias}.platform_order_id, '') || ' ' || + COALESCE(${alias}.product_name, '') || ' ' || + COALESCE(${alias}.material_json->'collect'->'fields'->>'gameId', '') || ' ' || + COALESCE(${alias}.material_json->'collect'->'fields'->>'gameNickname', '') || ' ' || + COALESCE(${alias}.material_json->'collect'->'fields'->>'gameAccount', '') || ' ' || + COALESCE(${alias}.material_json->'collect'->'fields'->>'roleName', '') || ' ' || + COALESCE(${alias}.material_json->'fields'->>'gameId', '') || ' ' || + COALESCE(${alias}.material_json->'fields'->>'gameNickname', '') || ' ' || + COALESCE(${alias}.material_json->>'gameId', '') || ' ' || + COALESCE(${alias}.material_json->>'gameNickname', '') + )` +} + +/** 打手订单筛选使用个人份额状态,而不是拼单母工单的 open 状态。 */ +function workerOrderStatusExpression(workerShareAlias = 'worker_share') { + return `CASE + WHEN wo.status IN ('accepted', 'cancelled', 'problem') THEN wo.status + ELSE COALESCE(CASE ${workerShareAlias}.status + WHEN 'joined' THEN 'in_progress' + WHEN 'submitted' THEN 'pending_acceptance' + WHEN 'accepted' THEN 'accepted' + WHEN 'cancelled' THEN 'cancelled' + ELSE wo.status + END, wo.status) + END` +} +export async function getWorkOrderByOrderItemId( + orderItemId: number | string, +): Promise { + const result = await query( + `${WORK_ORDER_SELECT} WHERE wo.order_item_id = $1 LIMIT 1`, + [Number(orderItemId)], + ) + return result.rows[0] || null +} + +export async function getWorkOrderById(workOrderId: number | string): Promise { + const result = await query(`${WORK_ORDER_SELECT} WHERE wo.id = $1 LIMIT 1`, [ + Number(workOrderId), + ]) + return result.rows[0] || null +} + +export async function findPendingMaterialWorkOrderByPlatformOrderId( + platformOrderId: string, +): Promise { + const result = await query( + `${WORK_ORDER_SELECT} + WHERE wo.platform_order_id = $1 AND wo.status = 'pending_material' + ORDER BY wo.id DESC + LIMIT 1`, + [String(platformOrderId || '').trim()], + ) + return result.rows[0] || null +} + +export async function listPendingMaterialWorkOrdersByPlatformOrderId( + platformOrderId: string, +): Promise { + const result = await query( + `${WORK_ORDER_SELECT} + WHERE wo.platform_order_id = $1 AND wo.status = 'pending_material' + ORDER BY wo.id ASC`, + [String(platformOrderId || '').trim()], + ) + return result.rows +} + +export async function listWorkOrdersByPlatformOrderId( + platformOrderId: string, +): Promise { + const result = await query( + `${WORK_ORDER_SELECT} + WHERE wo.platform_order_id = $1 + ORDER BY wo.id ASC`, + [String(platformOrderId || '').trim()], + ) + return result.rows +} + +export async function listWorkOrders({ + page = 1, + pageSize = 20, + statuses = [], + status = '', + workOrderId = 0, + keyword = '', + workerId = 0, + workerLevelIds = [], + categoryId = 0, + workerSharingId = 0, + hallCandidateLimit = 0, + visibleAfterIso = '', + excludeFilledSharing = false, + vipEvidencePending = false, + useOperationalStatus = false, + pinnedFirst = false, + sort = 'id_desc', +}: ListInput = {}): Promise<{ items: WorkOrderRow[]; total: number }> { + const effectiveStatuses = statuses.length > 0 ? statuses : status.trim() ? [status.trim()] : [] + const { joins, whereClause, params } = buildWorkOrderWhere({ + statuses: effectiveStatuses, + workOrderId, + keyword, + workerId, + workerLevelIds, + categoryId, + workerSharingId, + hallCandidateLimit, + visibleAfterIso, + excludeFilledSharing, + vipEvidencePending, + useOperationalStatus, + }) + const totalResult = await query<{ total: number }>( + ` + SELECT COUNT(*)::int AS total + FROM work_orders wo + ${joins} + ${whereClause} + `, + params, + ) + const offset = (page - 1) * pageSize + const orderBy = + sort === 'worker_claimed_at_desc' && workerSharingId + ? `COALESCE( + CASE + WHEN wo.assigned_worker_id = $1 THEN wo.assigned_at + END, + CASE + WHEN worker_share.status != 'cancelled' THEN worker_share.created_at + END + ) DESC NULLS LAST, wo.id DESC` + : 'wo.id DESC' + const effectiveOrderBy = pinnedFirst ? `wo.pinned_at DESC NULLS LAST, ${orderBy}` : orderBy + params.push(pageSize, offset) + const itemsResult = await query( + `${WORK_ORDER_SELECT} + ${joins} + ${whereClause} + ORDER BY ${effectiveOrderBy} + LIMIT $${params.length - 1} OFFSET $${params.length}`, + params, + ) + return { items: itemsResult.rows, total: Number(totalResult.rows[0]?.total || 0) } +} + +/** 按筛选条件汇总工单,统计范围不受当前分页影响。 */ +export async function getWorkOrderStatistics(input: ListInput = {}): Promise { + const { joins, whereClause, params } = buildWorkOrderWhere(input) + const vipEvidencePendingWhere = ` + wo.status = 'accepted' + AND wo.acceptance_json->>'reviewMode' = 'vip_auto' + AND wo.acceptance_json->>'evidenceStatus' = 'pending' + AND NULLIF(wo.acceptance_json->>'evidenceDueAt', '')::timestamptz > NOW() + ` + const operationalStatus = input.useOperationalStatus + ? adminOperationalStatusExpression('wo', 'operational_share_summary') + : 'wo.status' + const statusExpression = input.vipEvidencePending + ? `CASE WHEN ${vipEvidencePendingWhere} THEN 'vip_evidence_pending' ELSE ${operationalStatus} END` + : operationalStatus + const paymentExpression = ` + CASE + WHEN COALESCE(wo.material_json->'source'->>'orderAmountFen', '') ~ '^-?[0-9]+$' + THEN (wo.material_json->'source'->>'orderAmountFen')::bigint + ELSE COALESCE(o.total_amount, 0) + END + ` + const result = await query( + ` + SELECT + ${statusExpression} AS status, + COUNT(*)::int AS total_count, + COALESCE(SUM(${paymentExpression}), 0)::bigint AS user_payment_amount, + COALESCE(SUM(wo.reward_amount), 0)::bigint AS reward_amount, + COALESCE(SUM( + CASE + WHEN o.pay_status = 'refunded' OR o.order_status = 'refunded' + THEN ${paymentExpression} + ELSE 0 + END + ), 0)::bigint AS refund_amount + FROM work_orders wo + LEFT JOIN orders o ON o.id = wo.order_id + ${joins} + ${whereClause} + GROUP BY ${statusExpression} + `, + params, + ) + const statistics: WorkOrderStatistics = { + totalCount: 0, + userPaymentAmount: 0, + rewardAmount: 0, + refundAmount: 0, + statusCounts: {}, + } + for (const row of result.rows) { + const count = Number(row.total_count || 0) + statistics.totalCount += count + statistics.userPaymentAmount += Number(row.user_payment_amount || 0) + statistics.rewardAmount += Number(row.reward_amount || 0) + statistics.refundAmount += Number(row.refund_amount || 0) + statistics.statusCounts[row.status] = count + } + return statistics +} + +/** 汇总打手本人各订单状态的可结算金额,拼单按个人份额计算。 */ +export async function getWorkerOrderOverview(workerId: number | string): Promise< + Array<{ + status: string + orderCount: number + settlementAmount: number + vipEvidencePendingCount: number + vipEvidencePendingSettlementAmount: number + }> +> { + const normalizedWorkerId = Number(workerId) + const result = await query( + ` + WITH worker_orders AS ( + SELECT + CASE + WHEN wos.id IS NOT NULL AND wo.status NOT IN ('accepted', 'cancelled', 'problem') + THEN CASE wos.status + WHEN 'joined' THEN 'in_progress' + WHEN 'submitted' THEN 'pending_acceptance' + WHEN 'accepted' THEN 'accepted' + WHEN 'cancelled' THEN 'cancelled' + ELSE wo.status + END + ELSE wo.status + END AS status, + wo.acceptance_json, + wos.id IS NOT NULL AS is_worker_share, + CASE + WHEN wos.id IS NOT NULL THEN wos.share_reward + WHEN wo.assigned_worker_id = $1 THEN wo.reward_amount + ELSE 0 + END AS worker_settlement_amount + FROM work_orders wo + LEFT JOIN work_order_shares wos + ON wos.work_order_id = wo.id + AND wos.worker_id = $1 + AND wos.status != 'cancelled' + WHERE wo.assigned_worker_id = $1 + OR ( + wo.last_assigned_worker_id = $1 + AND ( + wo.status <> 'open' + OR EXISTS ( + SELECT 1 + FROM work_order_events timeout_event + WHERE timeout_event.work_order_id = wo.id + AND timeout_event.event_type = 'timeout_reopen' + ) + ) + ) + OR wos.id IS NOT NULL + ) + SELECT + status, + COUNT(*)::int AS order_count, + COALESCE( + SUM( + CASE + WHEN status IN ('in_progress', 'pending_acceptance', 'problem', 'accepted') + OR (status = 'open' AND is_worker_share) + THEN worker_settlement_amount + ELSE 0 + END + ), + 0 + )::int AS settlement_amount, + COUNT(*) FILTER ( + WHERE status = 'accepted' + AND acceptance_json->>'reviewMode' = 'vip_auto' + AND acceptance_json->>'evidenceStatus' = 'pending' + AND NULLIF(acceptance_json->>'evidenceDueAt', '')::timestamptz > NOW() + )::int AS vip_evidence_pending_count, + COALESCE( + SUM(worker_settlement_amount) FILTER ( + WHERE status = 'accepted' + AND acceptance_json->>'reviewMode' = 'vip_auto' + AND acceptance_json->>'evidenceStatus' = 'pending' + AND NULLIF(acceptance_json->>'evidenceDueAt', '')::timestamptz > NOW() + ), + 0 + )::int AS vip_evidence_pending_settlement_amount + FROM worker_orders + GROUP BY status + `, + [normalizedWorkerId], + ) + return result.rows.map((row) => ({ + status: row.status, + orderCount: Number(row.order_count || 0), + settlementAmount: Number(row.settlement_amount || 0), + vipEvidencePendingCount: Number(row.vip_evidence_pending_count || 0), + vipEvidencePendingSettlementAmount: Number(row.vip_evidence_pending_settlement_amount || 0), + })) +} +export async function countWorkerActiveOrders(workerId: number | string): Promise { + const result = await query<{ total: number }>( + ` + SELECT ( + SELECT COUNT(*)::int + FROM work_orders + WHERE assigned_worker_id = $1 + AND status = 'in_progress' + ) + ( + SELECT COUNT(*)::int + 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 = 'joined' + AND wo.status = 'open' + ) AS total + `, + [Number(workerId)], + ) + return Number(result.rows[0]?.total || 0) +} + +export async function countTimeoutEventsByWorkerIds( + workerIds: number[], +): Promise> { + const counts = new Map() + const uniqueIds = [...new Set(workerIds.map((id) => Number(id)).filter((id) => id > 0))] + if (uniqueIds.length === 0) return counts + const result = await query<{ worker_id: string; total: number }>( + ` + SELECT payload_json->>'workerId' AS worker_id, COUNT(*)::int AS total + FROM work_order_events + WHERE event_type LIKE 'timeout_%' + AND payload_json->>'workerId' IS NOT NULL + AND payload_json->>'workerId' != '' + AND payload_json->>'workerId' = ANY($1::text[]) + GROUP BY payload_json->>'workerId' + `, + [uniqueIds.map(String)], + ) + for (const row of result.rows) { + const workerId = Number(row.worker_id) + if (Number.isFinite(workerId) && workerId > 0) { + counts.set(workerId, Number(row.total || 0)) + } + } + return counts +} + +export async function countWorkerTimeoutEvents(workerId: number | string): Promise { + const counts = await countTimeoutEventsByWorkerIds([Number(workerId)]) + return counts.get(Number(workerId)) || 0 +} + +export async function getWorkOrderByIdWithClient( + client: PoolClient, + workOrderId: number | string, +): Promise { + const result = await client.query(`${WORK_ORDER_SELECT} WHERE wo.id = $1 LIMIT 1`, [ + Number(workOrderId), + ]) + return result.rows[0] || null +} + +export function buildWorkOrderWhere({ + statuses = [], + workOrderId = 0, + keyword = '', + workerId = 0, + workerLevelIds = [], + categoryId = 0, + workerSharingId = 0, + hallCandidateLimit = 0, + visibleAfterIso = '', + excludeFilledSharing = false, + vipEvidencePending = false, + useOperationalStatus = false, +}: ListInput) { + const filters: string[] = [] + const params: unknown[] = [] + const joins: string[] = [] + let workerSharePlaceholder = '' + + if (excludeFilledSharing) { + joins.push(` + LEFT JOIN ( + SELECT + wos.work_order_id, + SUM(wos.quantity) FILTER (WHERE wos.status != 'cancelled') AS active_share_quantity + FROM work_order_shares wos + GROUP BY wos.work_order_id + ) filter_share_summary ON filter_share_summary.work_order_id = wo.id`) + } + + // 先从打手的整单、历史单和拼单记录组成可见工单集合,再关联工单详情。 + // 这样 COUNT 和列表都不会先扫描全量工单,再逐行检查该打手的拼单状态。 + if (workerSharingId) { + params.push(workerSharingId) + workerSharePlaceholder = `$${params.length}` + joins.push(` + INNER JOIN ( + SELECT assigned.id AS work_order_id + FROM work_orders assigned + WHERE assigned.assigned_worker_id = ${workerSharePlaceholder} + UNION + SELECT historical.id AS work_order_id + FROM work_orders historical + WHERE historical.last_assigned_worker_id = ${workerSharePlaceholder} + AND ( + historical.status <> 'open' + OR EXISTS ( + SELECT 1 + FROM work_order_events timeout_event + WHERE timeout_event.work_order_id = historical.id + AND timeout_event.event_type = 'timeout_reopen' + ) + ) + UNION + SELECT share_source.work_order_id + FROM work_order_shares share_source + WHERE share_source.worker_id = ${workerSharePlaceholder} + AND share_source.status != 'cancelled' + ) worker_visible ON worker_visible.work_order_id = wo.id + LEFT JOIN work_order_shares worker_share + ON worker_share.work_order_id = wo.id + AND worker_share.worker_id = ${workerSharePlaceholder}`) + } + + if (useOperationalStatus) { + if (workerSharingId) { + // 已限定打手时只汇总其候选工单,利用 work_order_id 索引避免扫描整张拼单表。 + joins.push(` + LEFT JOIN LATERAL ( + SELECT + SUM(wos.quantity) FILTER (WHERE wos.status != 'cancelled') AS active_share_quantity, + COUNT(*) FILTER (WHERE wos.status = 'joined') AS pending_share_count + FROM work_order_shares wos + WHERE wos.work_order_id = wo.id + ) operational_share_summary ON TRUE`) + } else { + // 无打手范围时一次聚合全部拼单,避免对每笔工单重复执行关联子查询。 + joins.push(` + LEFT JOIN ( + SELECT + wos.work_order_id, + SUM(wos.quantity) FILTER (WHERE wos.status != 'cancelled') AS active_share_quantity, + COUNT(*) FILTER (WHERE wos.status = 'joined') AS pending_share_count + FROM work_order_shares wos + GROUP BY wos.work_order_id + ) operational_share_summary ON operational_share_summary.work_order_id = wo.id`) + } + } + const normalizedStatuses = [ + ...new Set(statuses.map((item) => String(item || '').trim()).filter(Boolean)), + ] + if (vipEvidencePending) { + const vipEvidencePendingWhere = ` + wo.status = 'accepted' + AND wo.acceptance_json->>'reviewMode' = 'vip_auto' + AND wo.acceptance_json->>'evidenceStatus' = 'pending' + AND NULLIF(wo.acceptance_json->>'evidenceDueAt', '')::timestamptz > NOW() + ` + if (normalizedStatuses.length > 0) { + params.push(normalizedStatuses) + filters.push( + `(${useOperationalStatus ? adminOperationalStatusExpression('wo', 'operational_share_summary') : 'wo.status'} = ANY($${params.length}::text[]) OR (${vipEvidencePendingWhere}))`, + ) + } else { + filters.push(`(${vipEvidencePendingWhere})`) + } + } else if (normalizedStatuses.length > 0) { + params.push(normalizedStatuses) + const statusesPlaceholder = `$${params.length}` + if (workerSharingId) { + filters.push(`${workerOrderStatusExpression()} = ANY(${statusesPlaceholder}::text[])`) + } else { + filters.push( + `${useOperationalStatus ? adminOperationalStatusExpression('wo', 'operational_share_summary') : 'wo.status'} = ANY(${statusesPlaceholder}::text[])`, + ) + } + } + if (workOrderId) { + params.push(workOrderId) + filters.push(`wo.id = $${params.length}`) + } + if (keyword) { + params.push(`%${keyword}%`) + filters.push(`${workOrderSearchExpression()} LIKE LOWER($${params.length})`) + } + if (workerId) { + params.push(workerId) + filters.push(`( + wo.assigned_worker_id = $${params.length} + OR EXISTS ( + SELECT 1 + FROM work_order_shares wos + WHERE wos.work_order_id = wo.id + AND wos.worker_id = $${params.length} + AND wos.status != 'cancelled' + ) + )`) + } + if (workerLevelIds.length > 0) { + params.push(workerLevelIds) + filters.push(`( + EXISTS ( + SELECT 1 + FROM worker_users wu + WHERE wu.id = wo.assigned_worker_id + AND wu.level_id = ANY($${params.length}::bigint[]) + ) + OR EXISTS ( + SELECT 1 + FROM work_order_shares wos + INNER JOIN worker_users wu ON wu.id = wos.worker_id + WHERE wos.work_order_id = wo.id + AND wos.status != 'cancelled' + AND wu.level_id = ANY($${params.length}::bigint[]) + ) + )`) + } + if (categoryId) { + params.push(categoryId) + filters.push(`wo.category_id = $${params.length}`) + } + if (hallCandidateLimit > 0) { + params.push(hallCandidateLimit) + filters.push(`wo.id IN (${buildHallCandidateIdsSql(`$${params.length}`)})`) + } + if (visibleAfterIso) { + params.push(visibleAfterIso) + filters.push(`(wo.hall_queued_at IS NULL OR wo.hall_queued_at <= $${params.length})`) + } + if (excludeFilledSharing) { + filters.push(`( + wo.sharing_enabled IS NOT TRUE + OR COALESCE(filter_share_summary.active_share_quantity, 0) < wo.sharing_total_quantity + )`) + } + return { + joins: joins.join('\n'), + whereClause: filters.length > 0 ? `WHERE ${filters.join(' AND ')}` : '', + params, + } +} + +/** 大厅容量先选出全局队首,分类、关键词和 VIP 延迟只能在该集合内继续筛选。 */ +function buildHallCandidateIdsSql(limitPlaceholder: string) { + return ` + SELECT candidate.id + FROM work_orders candidate + LEFT JOIN ( + SELECT + work_order_id, + SUM(quantity) FILTER (WHERE status != 'cancelled') AS active_share_quantity + FROM work_order_shares + GROUP BY work_order_id + ) candidate_shares ON candidate_shares.work_order_id = candidate.id + WHERE candidate.status = 'open' + AND ( + candidate.sharing_enabled IS NOT TRUE + OR COALESCE(candidate_shares.active_share_quantity, 0) < candidate.sharing_total_quantity + ) + ORDER BY candidate.hall_queued_at ASC NULLS LAST, candidate.id ASC + LIMIT ${limitPlaceholder} + ` +} + +/** 抢单操作在事务中复核容量与当前打手的可见延迟,避免通过订单 ID 绕过大厅队列。 */ +export async function isWorkOrderWithinHallCapacityWithClient( + client: PoolClient, + input: { workOrderId: number; hallCandidateLimit: number; visibleAfterIso?: string | undefined }, +) { + const result = await client.query<{ visible: boolean }>( + ` + SELECT EXISTS( + SELECT 1 + FROM work_orders wo + WHERE wo.id = $1 + AND wo.id IN (${buildHallCandidateIdsSql('$2')}) + AND ($3::timestamptz IS NULL OR wo.hall_queued_at IS NULL OR wo.hall_queued_at <= $3) + ) AS visible + `, + [ + input.workOrderId, + input.hallCandidateLimit, + input.visibleAfterIso ? input.visibleAfterIso : null, + ], + ) + return result.rows[0]?.visible === true +} 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 0b8264ea..e4bcffd3 100644 --- a/apps/backend/src/repositories/worker-platform/work-order-repo.ts +++ b/apps/backend/src/repositories/worker-platform/work-order-repo.ts @@ -9,6 +9,17 @@ import { resolveWorkOrderShareCancellationStatus, } from './shared.js' import { createWorkOrderEventWithClient } from './work-order-event-repo.js' +import { + countTimeoutEventsByWorkerIds, + countWorkerActiveOrders, + countWorkerTimeoutEvents, + getWorkOrderById, + getWorkOrderByIdWithClient, + getWorkOrderByOrderItemId, + getWorkerOrderOverview, + isWorkOrderWithinHallCapacityWithClient, + WORK_ORDER_SELECT, +} from './work-order-query-repo.js' import type { CreateWorkOrderInput, GrabWorkOrderResult, @@ -22,103 +33,7 @@ import type { import type { PoolClient } from 'pg' import { maybeUpgradeWorkerLevelWithClient } from './worker-ranking-repo.js' -export const WORK_ORDER_SELECT = ` - SELECT - wo.*, - COALESCE(share_summary.active_share_quantity, 0)::int AS active_share_quantity, - COALESCE(share_summary.pending_share_count, 0)::int AS pending_share_count, - wc.name AS category_name, - wu.username AS worker_username, - wu.display_name AS worker_display_name - FROM work_orders wo - LEFT JOIN work_categories wc ON wc.id = wo.category_id - LEFT JOIN worker_users wu ON wu.id = wo.assigned_worker_id - LEFT JOIN LATERAL ( - SELECT - SUM(wos.quantity) FILTER (WHERE wos.status != 'cancelled') AS active_share_quantity, - COUNT(*) FILTER (WHERE wos.status = 'joined') AS pending_share_count - FROM work_order_shares wos - WHERE wos.work_order_id = wo.id - ) share_summary ON TRUE -` - -type WorkerOrderOverviewRow = { - status: string - order_count: number - settlement_amount: number - vip_evidence_pending_count: number - vip_evidence_pending_settlement_amount: number -} - -type WorkOrderStatisticsRow = { - status: string - total_count: number - user_payment_amount: number | string - reward_amount: number | string - refund_amount: number | string -} - -/** 后台按母工单拼单参与情况计算业务状态,open 只代表大厅生命周期。 */ -function adminOperationalStatusExpression(alias = 'wo', shareSummaryAlias = '') { - if (shareSummaryAlias) { - return `CASE - WHEN ${alias}.status = 'open' - AND ${alias}.sharing_enabled IS TRUE - AND COALESCE(${shareSummaryAlias}.active_share_quantity, 0) > 0 - THEN CASE - WHEN COALESCE(${shareSummaryAlias}.active_share_quantity, 0) >= ${alias}.sharing_total_quantity - AND COALESCE(${shareSummaryAlias}.pending_share_count, 0) = 0 - THEN 'pending_acceptance' - ELSE 'in_progress' - END - ELSE ${alias}.status - END` - } - - return `CASE - WHEN ${alias}.status = 'open' - AND ${alias}.sharing_enabled IS TRUE - AND COALESCE((SELECT SUM(wos.quantity) FROM work_order_shares wos WHERE wos.work_order_id = ${alias}.id AND wos.status != 'cancelled'), 0) > 0 - THEN CASE - WHEN COALESCE((SELECT SUM(wos.quantity) FROM work_order_shares wos WHERE wos.work_order_id = ${alias}.id AND wos.status != 'cancelled'), 0) >= ${alias}.sharing_total_quantity - AND COALESCE((SELECT COUNT(*) FROM work_order_shares wos WHERE wos.work_order_id = ${alias}.id AND wos.status = 'joined'), 0) = 0 - THEN 'pending_acceptance' - ELSE 'in_progress' - END - ELSE ${alias}.status - END` -} - -/** 工单统一搜索字段,兼容当前资料结构和历史资料字段。 */ -function workOrderSearchExpression(alias = 'wo') { - return `LOWER( - COALESCE(${alias}.work_order_no, '') || ' ' || - COALESCE(${alias}.platform_order_id, '') || ' ' || - COALESCE(${alias}.product_name, '') || ' ' || - COALESCE(${alias}.material_json->'collect'->'fields'->>'gameId', '') || ' ' || - COALESCE(${alias}.material_json->'collect'->'fields'->>'gameNickname', '') || ' ' || - COALESCE(${alias}.material_json->'collect'->'fields'->>'gameAccount', '') || ' ' || - COALESCE(${alias}.material_json->'collect'->'fields'->>'roleName', '') || ' ' || - COALESCE(${alias}.material_json->'fields'->>'gameId', '') || ' ' || - COALESCE(${alias}.material_json->'fields'->>'gameNickname', '') || ' ' || - COALESCE(${alias}.material_json->>'gameId', '') || ' ' || - COALESCE(${alias}.material_json->>'gameNickname', '') - )` -} - -/** 打手订单筛选使用个人份额状态,而不是拼单母工单的 open 状态。 */ -function workerOrderStatusExpression(workerShareAlias = 'worker_share') { - return `CASE - WHEN wo.status IN ('accepted', 'cancelled', 'problem') THEN wo.status - ELSE COALESCE(CASE ${workerShareAlias}.status - WHEN 'joined' THEN 'in_progress' - WHEN 'submitted' THEN 'pending_acceptance' - WHEN 'accepted' THEN 'accepted' - WHEN 'cancelled' THEN 'cancelled' - ELSE wo.status - END, wo.status) - END` -} +export * from './work-order-query-repo.js' export async function createWorkOrder(input: CreateWorkOrderInput): Promise { const result = await query<{ id: number }>( @@ -162,286 +77,6 @@ export async function createWorkOrder(input: CreateWorkOrderInput): Promise { - const result = await query( - `${WORK_ORDER_SELECT} WHERE wo.order_item_id = $1 LIMIT 1`, - [Number(orderItemId)], - ) - return result.rows[0] || null -} - -export async function getWorkOrderById(workOrderId: number | string): Promise { - const result = await query(`${WORK_ORDER_SELECT} WHERE wo.id = $1 LIMIT 1`, [ - Number(workOrderId), - ]) - return result.rows[0] || null -} - -export async function findPendingMaterialWorkOrderByPlatformOrderId( - platformOrderId: string, -): Promise { - const result = await query( - `${WORK_ORDER_SELECT} - WHERE wo.platform_order_id = $1 AND wo.status = 'pending_material' - ORDER BY wo.id DESC - LIMIT 1`, - [String(platformOrderId || '').trim()], - ) - return result.rows[0] || null -} - -export async function listPendingMaterialWorkOrdersByPlatformOrderId( - platformOrderId: string, -): Promise { - const result = await query( - `${WORK_ORDER_SELECT} - WHERE wo.platform_order_id = $1 AND wo.status = 'pending_material' - ORDER BY wo.id ASC`, - [String(platformOrderId || '').trim()], - ) - return result.rows -} - -export async function listWorkOrdersByPlatformOrderId( - platformOrderId: string, -): Promise { - const result = await query( - `${WORK_ORDER_SELECT} - WHERE wo.platform_order_id = $1 - ORDER BY wo.id ASC`, - [String(platformOrderId || '').trim()], - ) - return result.rows -} - -export async function listWorkOrders({ - page = 1, - pageSize = 20, - statuses = [], - status = '', - workOrderId = 0, - keyword = '', - workerId = 0, - workerLevelIds = [], - categoryId = 0, - workerSharingId = 0, - hallCandidateLimit = 0, - visibleAfterIso = '', - excludeFilledSharing = false, - vipEvidencePending = false, - useOperationalStatus = false, - pinnedFirst = false, - sort = 'id_desc', -}: ListInput = {}): Promise<{ items: WorkOrderRow[]; total: number }> { - const effectiveStatuses = statuses.length > 0 ? statuses : status.trim() ? [status.trim()] : [] - const { joins, whereClause, params } = buildWorkOrderWhere({ - statuses: effectiveStatuses, - workOrderId, - keyword, - workerId, - workerLevelIds, - categoryId, - workerSharingId, - hallCandidateLimit, - visibleAfterIso, - excludeFilledSharing, - vipEvidencePending, - useOperationalStatus, - }) - const totalResult = await query<{ total: number }>( - ` - SELECT COUNT(*)::int AS total - FROM work_orders wo - ${joins} - ${whereClause} - `, - params, - ) - const offset = (page - 1) * pageSize - const orderBy = - sort === 'worker_claimed_at_desc' && workerSharingId - ? `COALESCE( - CASE - WHEN wo.assigned_worker_id = $1 THEN wo.assigned_at - END, - CASE - WHEN worker_share.status != 'cancelled' THEN worker_share.created_at - END - ) DESC NULLS LAST, wo.id DESC` - : 'wo.id DESC' - const effectiveOrderBy = pinnedFirst ? `wo.pinned_at DESC NULLS LAST, ${orderBy}` : orderBy - params.push(pageSize, offset) - const itemsResult = await query( - `${WORK_ORDER_SELECT} - ${joins} - ${whereClause} - ORDER BY ${effectiveOrderBy} - LIMIT $${params.length - 1} OFFSET $${params.length}`, - params, - ) - return { items: itemsResult.rows, total: Number(totalResult.rows[0]?.total || 0) } -} - -/** 按筛选条件汇总工单,统计范围不受当前分页影响。 */ -export async function getWorkOrderStatistics(input: ListInput = {}): Promise { - const { joins, whereClause, params } = buildWorkOrderWhere(input) - const vipEvidencePendingWhere = ` - wo.status = 'accepted' - AND wo.acceptance_json->>'reviewMode' = 'vip_auto' - AND wo.acceptance_json->>'evidenceStatus' = 'pending' - AND NULLIF(wo.acceptance_json->>'evidenceDueAt', '')::timestamptz > NOW() - ` - const operationalStatus = input.useOperationalStatus - ? adminOperationalStatusExpression('wo', 'operational_share_summary') - : 'wo.status' - const statusExpression = input.vipEvidencePending - ? `CASE WHEN ${vipEvidencePendingWhere} THEN 'vip_evidence_pending' ELSE ${operationalStatus} END` - : operationalStatus - const paymentExpression = ` - CASE - WHEN COALESCE(wo.material_json->'source'->>'orderAmountFen', '') ~ '^-?[0-9]+$' - THEN (wo.material_json->'source'->>'orderAmountFen')::bigint - ELSE COALESCE(o.total_amount, 0) - END - ` - const result = await query( - ` - SELECT - ${statusExpression} AS status, - COUNT(*)::int AS total_count, - COALESCE(SUM(${paymentExpression}), 0)::bigint AS user_payment_amount, - COALESCE(SUM(wo.reward_amount), 0)::bigint AS reward_amount, - COALESCE(SUM( - CASE - WHEN o.pay_status = 'refunded' OR o.order_status = 'refunded' - THEN ${paymentExpression} - ELSE 0 - END - ), 0)::bigint AS refund_amount - FROM work_orders wo - LEFT JOIN orders o ON o.id = wo.order_id - ${joins} - ${whereClause} - GROUP BY ${statusExpression} - `, - params, - ) - const statistics: WorkOrderStatistics = { - totalCount: 0, - userPaymentAmount: 0, - rewardAmount: 0, - refundAmount: 0, - statusCounts: {}, - } - for (const row of result.rows) { - const count = Number(row.total_count || 0) - statistics.totalCount += count - statistics.userPaymentAmount += Number(row.user_payment_amount || 0) - statistics.rewardAmount += Number(row.reward_amount || 0) - statistics.refundAmount += Number(row.refund_amount || 0) - statistics.statusCounts[row.status] = count - } - return statistics -} - -/** 汇总打手本人各订单状态的可结算金额,拼单按个人份额计算。 */ -export async function getWorkerOrderOverview(workerId: number | string): Promise< - Array<{ - status: string - orderCount: number - settlementAmount: number - vipEvidencePendingCount: number - vipEvidencePendingSettlementAmount: number - }> -> { - const normalizedWorkerId = Number(workerId) - const result = await query( - ` - WITH worker_orders AS ( - SELECT - CASE - WHEN wos.id IS NOT NULL AND wo.status NOT IN ('accepted', 'cancelled', 'problem') - THEN CASE wos.status - WHEN 'joined' THEN 'in_progress' - WHEN 'submitted' THEN 'pending_acceptance' - WHEN 'accepted' THEN 'accepted' - WHEN 'cancelled' THEN 'cancelled' - ELSE wo.status - END - ELSE wo.status - END AS status, - wo.acceptance_json, - wos.id IS NOT NULL AS is_worker_share, - CASE - WHEN wos.id IS NOT NULL THEN wos.share_reward - WHEN wo.assigned_worker_id = $1 THEN wo.reward_amount - ELSE 0 - END AS worker_settlement_amount - FROM work_orders wo - LEFT JOIN work_order_shares wos - ON wos.work_order_id = wo.id - AND wos.worker_id = $1 - AND wos.status != 'cancelled' - WHERE wo.assigned_worker_id = $1 - OR ( - wo.last_assigned_worker_id = $1 - AND ( - wo.status <> 'open' - OR EXISTS ( - SELECT 1 - FROM work_order_events timeout_event - WHERE timeout_event.work_order_id = wo.id - AND timeout_event.event_type = 'timeout_reopen' - ) - ) - ) - OR wos.id IS NOT NULL - ) - SELECT - status, - COUNT(*)::int AS order_count, - COALESCE( - SUM( - CASE - WHEN status IN ('in_progress', 'pending_acceptance', 'problem', 'accepted') - OR (status = 'open' AND is_worker_share) - THEN worker_settlement_amount - ELSE 0 - END - ), - 0 - )::int AS settlement_amount, - COUNT(*) FILTER ( - WHERE status = 'accepted' - AND acceptance_json->>'reviewMode' = 'vip_auto' - AND acceptance_json->>'evidenceStatus' = 'pending' - AND NULLIF(acceptance_json->>'evidenceDueAt', '')::timestamptz > NOW() - )::int AS vip_evidence_pending_count, - COALESCE( - SUM(worker_settlement_amount) FILTER ( - WHERE status = 'accepted' - AND acceptance_json->>'reviewMode' = 'vip_auto' - AND acceptance_json->>'evidenceStatus' = 'pending' - AND NULLIF(acceptance_json->>'evidenceDueAt', '')::timestamptz > NOW() - ), - 0 - )::int AS vip_evidence_pending_settlement_amount - FROM worker_orders - GROUP BY status - `, - [normalizedWorkerId], - ) - return result.rows.map((row) => ({ - status: row.status, - orderCount: Number(row.order_count || 0), - settlementAmount: Number(row.settlement_amount || 0), - vipEvidencePendingCount: Number(row.vip_evidence_pending_count || 0), - vipEvidencePendingSettlementAmount: Number(row.vip_evidence_pending_settlement_amount || 0), - })) -} - export async function updateWorkOrder( workOrderId: number | string, patch: Partial< @@ -2491,58 +2126,22 @@ async function enqueueDepositUnfreezeWithClient( ) } -export async function countWorkerActiveOrders(workerId: number | string): Promise { - const result = await query<{ total: number }>( - ` - SELECT ( - SELECT COUNT(*)::int - FROM work_orders - WHERE assigned_worker_id = $1 - AND status = 'in_progress' - ) + ( - SELECT COUNT(*)::int - 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 = 'joined' - AND wo.status = 'open' - ) AS total - `, - [Number(workerId)], +export function resolveOutstandingDepositAmount( + totals: + | { + frozen_amount?: number + released_amount?: number + deducted_amount?: number + } + | null + | undefined, +) { + return Math.max( + 0, + Number(totals?.frozen_amount || 0) - + Number(totals?.released_amount || 0) - + Number(totals?.deducted_amount || 0), ) - return Number(result.rows[0]?.total || 0) -} - -export async function countTimeoutEventsByWorkerIds( - workerIds: number[], -): Promise> { - const counts = new Map() - const uniqueIds = [...new Set(workerIds.map((id) => Number(id)).filter((id) => id > 0))] - if (uniqueIds.length === 0) return counts - const result = await query<{ worker_id: string; total: number }>( - ` - SELECT payload_json->>'workerId' AS worker_id, COUNT(*)::int AS total - FROM work_order_events - WHERE event_type LIKE 'timeout_%' - AND payload_json->>'workerId' IS NOT NULL - AND payload_json->>'workerId' != '' - AND payload_json->>'workerId' = ANY($1::text[]) - GROUP BY payload_json->>'workerId' - `, - [uniqueIds.map(String)], - ) - for (const row of result.rows) { - const workerId = Number(row.worker_id) - if (Number.isFinite(workerId) && workerId > 0) { - counts.set(workerId, Number(row.total || 0)) - } - } - return counts -} - -export async function countWorkerTimeoutEvents(workerId: number | string): Promise { - const counts = await countTimeoutEventsByWorkerIds([Number(workerId)]) - return counts.get(Number(workerId)) || 0 } export async function getOutstandingDepositAmountWithClient( @@ -2593,259 +2192,3 @@ export async function countWorkerActiveOrdersWithClient( ) return Number(result.rows[0]?.total || 0) } - -export function resolveOutstandingDepositAmount( - totals: - | { - frozen_amount?: number - released_amount?: number - deducted_amount?: number - } - | null - | undefined, -) { - return Math.max( - 0, - Number(totals?.frozen_amount || 0) - - Number(totals?.released_amount || 0) - - Number(totals?.deducted_amount || 0), - ) -} - -export async function getWorkOrderByIdWithClient( - client: PoolClient, - workOrderId: number | string, -): Promise { - const result = await client.query(`${WORK_ORDER_SELECT} WHERE wo.id = $1 LIMIT 1`, [ - Number(workOrderId), - ]) - return result.rows[0] || null -} - -export function buildWorkOrderWhere({ - statuses = [], - workOrderId = 0, - keyword = '', - workerId = 0, - workerLevelIds = [], - categoryId = 0, - workerSharingId = 0, - hallCandidateLimit = 0, - visibleAfterIso = '', - excludeFilledSharing = false, - vipEvidencePending = false, - useOperationalStatus = false, -}: ListInput) { - const filters: string[] = [] - const params: unknown[] = [] - const joins: string[] = [] - let workerSharePlaceholder = '' - - if (excludeFilledSharing) { - joins.push(` - LEFT JOIN ( - SELECT - wos.work_order_id, - SUM(wos.quantity) FILTER (WHERE wos.status != 'cancelled') AS active_share_quantity - FROM work_order_shares wos - GROUP BY wos.work_order_id - ) filter_share_summary ON filter_share_summary.work_order_id = wo.id`) - } - - // 先从打手的整单、历史单和拼单记录组成可见工单集合,再关联工单详情。 - // 这样 COUNT 和列表都不会先扫描全量工单,再逐行检查该打手的拼单状态。 - if (workerSharingId) { - params.push(workerSharingId) - workerSharePlaceholder = `$${params.length}` - joins.push(` - INNER JOIN ( - SELECT assigned.id AS work_order_id - FROM work_orders assigned - WHERE assigned.assigned_worker_id = ${workerSharePlaceholder} - UNION - SELECT historical.id AS work_order_id - FROM work_orders historical - WHERE historical.last_assigned_worker_id = ${workerSharePlaceholder} - AND ( - historical.status <> 'open' - OR EXISTS ( - SELECT 1 - FROM work_order_events timeout_event - WHERE timeout_event.work_order_id = historical.id - AND timeout_event.event_type = 'timeout_reopen' - ) - ) - UNION - SELECT share_source.work_order_id - FROM work_order_shares share_source - WHERE share_source.worker_id = ${workerSharePlaceholder} - AND share_source.status != 'cancelled' - ) worker_visible ON worker_visible.work_order_id = wo.id - LEFT JOIN work_order_shares worker_share - ON worker_share.work_order_id = wo.id - AND worker_share.worker_id = ${workerSharePlaceholder}`) - } - - if (useOperationalStatus) { - if (workerSharingId) { - // 已限定打手时只汇总其候选工单,利用 work_order_id 索引避免扫描整张拼单表。 - joins.push(` - LEFT JOIN LATERAL ( - SELECT - SUM(wos.quantity) FILTER (WHERE wos.status != 'cancelled') AS active_share_quantity, - COUNT(*) FILTER (WHERE wos.status = 'joined') AS pending_share_count - FROM work_order_shares wos - WHERE wos.work_order_id = wo.id - ) operational_share_summary ON TRUE`) - } else { - // 无打手范围时一次聚合全部拼单,避免对每笔工单重复执行关联子查询。 - joins.push(` - LEFT JOIN ( - SELECT - wos.work_order_id, - SUM(wos.quantity) FILTER (WHERE wos.status != 'cancelled') AS active_share_quantity, - COUNT(*) FILTER (WHERE wos.status = 'joined') AS pending_share_count - FROM work_order_shares wos - GROUP BY wos.work_order_id - ) operational_share_summary ON operational_share_summary.work_order_id = wo.id`) - } - } - const normalizedStatuses = [ - ...new Set(statuses.map((item) => String(item || '').trim()).filter(Boolean)), - ] - if (vipEvidencePending) { - const vipEvidencePendingWhere = ` - wo.status = 'accepted' - AND wo.acceptance_json->>'reviewMode' = 'vip_auto' - AND wo.acceptance_json->>'evidenceStatus' = 'pending' - AND NULLIF(wo.acceptance_json->>'evidenceDueAt', '')::timestamptz > NOW() - ` - if (normalizedStatuses.length > 0) { - params.push(normalizedStatuses) - filters.push( - `(${useOperationalStatus ? adminOperationalStatusExpression('wo', 'operational_share_summary') : 'wo.status'} = ANY($${params.length}::text[]) OR (${vipEvidencePendingWhere}))`, - ) - } else { - filters.push(`(${vipEvidencePendingWhere})`) - } - } else if (normalizedStatuses.length > 0) { - params.push(normalizedStatuses) - const statusesPlaceholder = `$${params.length}` - if (workerSharingId) { - filters.push(`${workerOrderStatusExpression()} = ANY(${statusesPlaceholder}::text[])`) - } else { - filters.push( - `${useOperationalStatus ? adminOperationalStatusExpression('wo', 'operational_share_summary') : 'wo.status'} = ANY(${statusesPlaceholder}::text[])`, - ) - } - } - if (workOrderId) { - params.push(workOrderId) - filters.push(`wo.id = $${params.length}`) - } - if (keyword) { - params.push(`%${keyword}%`) - filters.push(`${workOrderSearchExpression()} LIKE LOWER($${params.length})`) - } - if (workerId) { - params.push(workerId) - filters.push(`( - wo.assigned_worker_id = $${params.length} - OR EXISTS ( - SELECT 1 - FROM work_order_shares wos - WHERE wos.work_order_id = wo.id - AND wos.worker_id = $${params.length} - AND wos.status != 'cancelled' - ) - )`) - } - if (workerLevelIds.length > 0) { - params.push(workerLevelIds) - filters.push(`( - EXISTS ( - SELECT 1 - FROM worker_users wu - WHERE wu.id = wo.assigned_worker_id - AND wu.level_id = ANY($${params.length}::bigint[]) - ) - OR EXISTS ( - SELECT 1 - FROM work_order_shares wos - INNER JOIN worker_users wu ON wu.id = wos.worker_id - WHERE wos.work_order_id = wo.id - AND wos.status != 'cancelled' - AND wu.level_id = ANY($${params.length}::bigint[]) - ) - )`) - } - if (categoryId) { - params.push(categoryId) - filters.push(`wo.category_id = $${params.length}`) - } - if (hallCandidateLimit > 0) { - params.push(hallCandidateLimit) - filters.push(`wo.id IN (${buildHallCandidateIdsSql(`$${params.length}`)})`) - } - if (visibleAfterIso) { - params.push(visibleAfterIso) - filters.push(`(wo.hall_queued_at IS NULL OR wo.hall_queued_at <= $${params.length})`) - } - if (excludeFilledSharing) { - filters.push(`( - wo.sharing_enabled IS NOT TRUE - OR COALESCE(filter_share_summary.active_share_quantity, 0) < wo.sharing_total_quantity - )`) - } - return { - joins: joins.join('\n'), - whereClause: filters.length > 0 ? `WHERE ${filters.join(' AND ')}` : '', - params, - } -} - -/** 大厅容量先选出全局队首,分类、关键词和 VIP 延迟只能在该集合内继续筛选。 */ -function buildHallCandidateIdsSql(limitPlaceholder: string) { - return ` - SELECT candidate.id - FROM work_orders candidate - LEFT JOIN ( - SELECT - work_order_id, - SUM(quantity) FILTER (WHERE status != 'cancelled') AS active_share_quantity - FROM work_order_shares - GROUP BY work_order_id - ) candidate_shares ON candidate_shares.work_order_id = candidate.id - WHERE candidate.status = 'open' - AND ( - candidate.sharing_enabled IS NOT TRUE - OR COALESCE(candidate_shares.active_share_quantity, 0) < candidate.sharing_total_quantity - ) - ORDER BY candidate.hall_queued_at ASC NULLS LAST, candidate.id ASC - LIMIT ${limitPlaceholder} - ` -} - -/** 抢单操作在事务中复核容量与当前打手的可见延迟,避免通过订单 ID 绕过大厅队列。 */ -export async function isWorkOrderWithinHallCapacityWithClient( - client: PoolClient, - input: { workOrderId: number; hallCandidateLimit: number; visibleAfterIso?: string | undefined }, -) { - const result = await client.query<{ visible: boolean }>( - ` - SELECT EXISTS( - SELECT 1 - FROM work_orders wo - WHERE wo.id = $1 - AND wo.id IN (${buildHallCandidateIdsSql('$2')}) - AND ($3::timestamptz IS NULL OR wo.hall_queued_at IS NULL OR wo.hall_queued_at <= $3) - ) AS visible - `, - [ - input.workOrderId, - input.hallCandidateLimit, - input.visibleAfterIso ? input.visibleAfterIso : null, - ], - ) - return result.rows[0]?.visible === true -}