diff --git a/apps/backend/src/db/postgres.integration.test.ts b/apps/backend/src/db/postgres.integration.test.ts index c9d28eeb..28c492fc 100644 --- a/apps/backend/src/db/postgres.integration.test.ts +++ b/apps/backend/src/db/postgres.integration.test.ts @@ -35,7 +35,7 @@ test( const tableName = `integration_webhook_${Date.now()}` try { await pool.query( - `CREATE TEMP TABLE ${tableName} ( + `CREATE TABLE ${tableName} ( id serial primary key, provider text not null, event_id text not null, @@ -54,6 +54,7 @@ test( const result = await pool.query(`SELECT COUNT(*)::int AS count FROM ${tableName}`) assert.equal(result.rows[0].count, 1) } finally { + await pool.query(`DROP TABLE IF EXISTS ${tableName}`) await pool.end() } }, 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 index 15a19d68..b84151c2 100644 --- a/apps/backend/src/repositories/worker-platform/work-order-query-repo.ts +++ b/apps/backend/src/repositories/worker-platform/work-order-query-repo.ts @@ -690,6 +690,7 @@ export function buildWorkOrderWhere({ const filters: string[] = [] const params: unknown[] = [] const joins: string[] = [] + const normalizedKeyword = String(keyword || '').trim() if (createdFrom) { params.push(createdFrom) filters.push( @@ -808,8 +809,8 @@ export function buildWorkOrderWhere({ params.push(workOrderId) filters.push(`wo.id = $${params.length}`) } - if (keyword) { - params.push(`%${keyword}%`) + if (normalizedKeyword) { + params.push(`%${normalizedKeyword}%`) filters.push(`${workOrderSearchExpression()} LIKE LOWER($${params.length})`) } if (workerId) { diff --git a/test-fix-artifacts/BASELINE_postgres_integration.test.ts b/test-fix-artifacts/BASELINE_postgres_integration.test.ts new file mode 100644 index 00000000..c9d28eeb --- /dev/null +++ b/test-fix-artifacts/BASELINE_postgres_integration.test.ts @@ -0,0 +1,60 @@ +import assert from 'node:assert/strict' +import test from 'node:test' +import { Pool } from 'pg' + +const connectionString = process.env.TEST_POSTGRES_URL + +test('real PostgreSQL transaction rolls back atomically', { skip: !connectionString }, async () => { + const pool = new Pool({ connectionString }) + const client = await pool.connect() + const tableName = `integration_tx_${Date.now()}` + try { + await client.query(`CREATE TEMP TABLE ${tableName} (id integer primary key, value text)`) + await client.query('BEGIN') + await client.query(`INSERT INTO ${tableName} (id, value) VALUES (1, 'before-error')`) + try { + await client.query(`INSERT INTO ${tableName} (id, value) VALUES (1, 'duplicate')`) + } catch { + await client.query('ROLLBACK') + } + const result = await client.query(`SELECT COUNT(*)::int AS count FROM ${tableName}`) + assert.equal(result.rows[0].count, 0) + } finally { + client.release() + await pool.end() + } +}) + +test( + 'real PostgreSQL enforces callback replay idempotency under concurrency', + { + skip: !connectionString, + }, + async () => { + const pool = new Pool({ connectionString }) + const tableName = `integration_webhook_${Date.now()}` + try { + await pool.query( + `CREATE TEMP TABLE ${tableName} ( + id serial primary key, + provider text not null, + event_id text not null, + unique(provider, event_id) + )`, + ) + await Promise.all( + Array.from({ length: 8 }, () => + pool.query( + `INSERT INTO ${tableName} (provider, event_id) + VALUES ('integration', 'replay-1') + ON CONFLICT (provider, event_id) DO NOTHING`, + ), + ), + ) + const result = await pool.query(`SELECT COUNT(*)::int AS count FROM ${tableName}`) + assert.equal(result.rows[0].count, 1) + } finally { + await pool.end() + } + }, +) diff --git a/test-fix-artifacts/BASELINE_worker_query.ts b/test-fix-artifacts/BASELINE_worker_query.ts new file mode 100644 index 00000000..15a19d68 --- /dev/null +++ b/test-fix-artifacts/BASELINE_worker_query.ts @@ -0,0 +1,960 @@ +import { query } from '../../db/client.js' +import type { ListInput, WorkOrderRow, WorkOrderStatistics } from './types.js' +import type { PoolClient } from 'pg' +import { + normalizeWorkOrderAcceptanceMode, + WORK_ORDER_GIFT_PHASE, +} from '../../domain/work-order-acceptance-mode.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, + cancel_info.cancel_reason, + cancel_info.cancel_event_type, + cancel_info.cancel_actor_type, + cancel_info.cancel_actor_id, + cancel_info.cancelled_at + 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 + ev.event_type AS cancel_event_type, + ev.actor_type AS cancel_actor_type, + ev.actor_id AS cancel_actor_id, + ev.created_at AS cancelled_at, + BTRIM(COALESCE( + NULLIF(BTRIM(COALESCE(ev.payload_json->>'reason', '')), ''), + NULLIF(BTRIM(COALESCE(ev.payload_json->>'note', '')), ''), + CASE ev.event_type + WHEN 'timeout_cancel_release' THEN '任务超时未完成,系统自动取消并退还押金' + WHEN 'timeout_cancel_deduct' THEN '任务超时未完成,系统自动取消并扣除押金' + ELSE '' + END + )) AS cancel_reason + FROM work_order_events ev + WHERE ev.work_order_id = wo.id + AND wo.status = 'cancelled' + AND ev.event_type IN ( + 'cancelled_by_admin', + 'problem_cancel_release', + 'problem_cancel_deduct', + 'timeout_cancel_release', + 'timeout_cancel_deduct' + ) + ORDER BY ev.id DESC + LIMIT 1 + ) cancel_info ON TRUE + 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 +` + +// Keep the page and total in one filtered query. This avoids running the same +// joins and share aggregation twice for every list request. +const LIST_WORK_ORDER_SELECT = WORK_ORDER_SELECT.replace( + ' SELECT\n', + ' SELECT\n COUNT(*) OVER()::int AS total_count,\n', +) + +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 +} + +type TimeoutCountCacheEntry = { + expiresAt: number + counts: Map +} + +// 后台打手列表会在短时间内重复请求同一页;超时次数只是展示字段,允许短 TTL。 +const TIMEOUT_COUNT_CACHE_TTL_MS = 5_000 +const TIMEOUT_COUNT_CACHE_MAX_ENTRIES = 256 +const timeoutCountCache = new Map() + +/** 后台按母工单拼单参与情况计算业务状态,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', '') + )` +} + +/** + * 打手订单个人状态:拼单按个人份额状态计算。 + * 传入 workerId 时,代练中但存在待审核取消申请或未处理反馈的订单归入问题单, + * 客服驳回/解决后自动回到代练中(计算型归类,无状态迁移)。 + */ +export function workerOrderStatusExpression(workerId = 0, workerShareAlias = 'worker_share') { + const baseExpression = `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` + const normalizedWorkerId = Math.max(0, Math.floor(Number(workerId) || 0)) + if (normalizedWorkerId <= 0) return baseExpression + const pendingHandlingExists = `( + EXISTS ( + SELECT 1 + FROM worker_order_cancel_requests wcr + WHERE wcr.work_order_id = wo.id + AND wcr.worker_id = ${normalizedWorkerId} + AND wcr.status = 'pending' + ) OR EXISTS ( + SELECT 1 + FROM worker_order_feedbacks wf + WHERE wf.work_order_id = wo.id + AND wf.worker_id = ${normalizedWorkerId} + AND wf.status = 'pending' + ) + )` + return `CASE + WHEN (${baseExpression}) = 'in_progress' AND ${pendingHandlingExists} THEN 'problem' + ELSE ${baseExpression} + 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, + hallCategoryLimits = {}, + visibleAfterIso = '', + excludeFilledSharing = false, + vipEvidencePending = false, + acceptanceMode = '', + giftPhase = '', + useOperationalStatus = false, + pinnedFirst = false, + sort = 'id_desc', + createdFrom = '', + createdTo = '', +}: 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, + hallCategoryLimits, + visibleAfterIso, + excludeFilledSharing, + vipEvidencePending, + acceptanceMode, + giftPhase, + useOperationalStatus, + createdFrom, + createdTo, + }) + 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( + `${LIST_WORK_ORDER_SELECT} + ${joins} + ${whereClause} + ORDER BY ${effectiveOrderBy} + LIMIT $${params.length - 1} OFFSET $${params.length}`, + params, + ) + // A page beyond the end has no row from which to read the window count. + // Keep that uncommon pagination edge case correct without restoring a + // second query for normal requests. + let total = Number( + (itemsResult.rows[0] as WorkOrderRow & { total_count?: number })?.total_count || 0, + ) + if (itemsResult.rows.length === 0 && page > 1) { + const totalResult = await query<{ total: number }>( + `SELECT COUNT(*)::int AS total FROM work_orders wo ${joins} ${whereClause}`, + params.slice(0, -2), + ) + total = Number(totalResult.rows[0]?.total || 0) + } + return { + items: itemsResult.rows, + total, + } +} + +/** 按筛选条件汇总工单,统计范围不受当前分页影响。 */ +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 visible_orders AS ( + SELECT wo.id AS work_order_id + FROM work_orders wo + WHERE wo.assigned_worker_id = $1 + UNION + SELECT wo.id AS work_order_id + FROM work_orders wo + WHERE 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' + ) + ) + UNION + SELECT wos.work_order_id + FROM work_order_shares wos + WHERE wos.worker_id = $1 + AND wos.status != 'cancelled' + ), + worker_base AS ( + SELECT + wo.id, + wo.status AS wo_status, + wos.status AS wos_status, + 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, + wo.acceptance_json + FROM work_orders wo + INNER JOIN visible_orders vo ON vo.work_order_id = wo.id + LEFT JOIN work_order_shares wos + ON wos.work_order_id = wo.id + AND wos.worker_id = $1 + AND wos.status != 'cancelled' + ), + worker_orders AS ( + SELECT + CASE + WHEN ( + CASE + WHEN is_worker_share 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 + ) = 'in_progress' AND ( + EXISTS ( + SELECT 1 + FROM worker_order_cancel_requests wcr + WHERE wcr.work_order_id = worker_base.id + AND wcr.worker_id = $1 + AND wcr.status = 'pending' + ) OR EXISTS ( + SELECT 1 + FROM worker_order_feedbacks wf + WHERE wf.work_order_id = worker_base.id + AND wf.worker_id = $1 + AND wf.status = 'pending' + ) + ) THEN 'problem' + ELSE CASE + WHEN is_worker_share 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 + END AS status, + acceptance_json, + is_worker_share, + worker_settlement_amount + FROM worker_base + ) + 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 listOverdueWorkOrders({ + limit = 50, + workerId = 0, +}: { limit?: number; workerId?: number } = {}): Promise { + const params: unknown[] = [] + let workerClause = '' + if (workerId > 0) { + params.push(workerId) + workerClause = `AND wo.assigned_worker_id = $${params.length}` + } + params.push(limit) + const result = await query( + `${WORK_ORDER_SELECT} + WHERE wo.status = 'in_progress' + AND wo.deadline_at IS NOT NULL + AND wo.deadline_at < NOW() + AND NOT (wo.acceptance_mode = 'friend_gift' AND wo.gift_phase = 'material_submitted') + ${workerClause} + ORDER BY wo.deadline_at ASC + LIMIT $${params.length}`, + params, + ) + return result.rows +} + +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 countWorkerActiveOrdersWithClient( + client: PoolClient, + workerId: number, +): Promise { + const result = await client.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 + `, + [workerId], + ) + return Number(result.rows[0]?.total || 0) +} + +export async function countTimeoutEventsByWorkerIds( + workerIds: number[], +): Promise> { + const uniqueIds = [...new Set(workerIds.map((id) => Number(id)).filter((id) => id > 0))] + if (uniqueIds.length === 0) return new Map() + uniqueIds.sort((left, right) => left - right) + const cacheKey = uniqueIds.join(',') + const now = Date.now() + const cached = timeoutCountCache.get(cacheKey) + if (cached && cached.expiresAt > now) { + return new Map(cached.counts) + } + + const counts = new Map() + 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)) + } + } + if (timeoutCountCache.size >= TIMEOUT_COUNT_CACHE_MAX_ENTRIES) { + const oldestKey = timeoutCountCache.keys().next().value + if (oldestKey) timeoutCountCache.delete(oldestKey) + } + timeoutCountCache.set(cacheKey, { + expiresAt: now + TIMEOUT_COUNT_CACHE_TTL_MS, + counts: new Map(counts), + }) + 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 countWorkerCancellationsSince( + workerId: number | string, + sinceIso: string, +): Promise { + const result = await query<{ total: number }>( + ` + SELECT COUNT(*)::int AS total + FROM work_order_events + WHERE actor_type = 'worker' + AND actor_id = $1 + AND event_type = 'cancelled_by_worker' + AND created_at >= $2 + `, + [String(workerId), sinceIso], + ) + return Number(result.rows[0]?.total || 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, + hallCategoryLimits = {}, + visibleAfterIso = '', + excludeFilledSharing = false, + vipEvidencePending = false, + acceptanceMode = '', + giftPhase = '', + useOperationalStatus = false, + createdFrom = '', + createdTo = '', +}: ListInput) { + const filters: string[] = [] + const params: unknown[] = [] + const joins: string[] = [] + if (createdFrom) { + params.push(createdFrom) + filters.push( + `wo.created_at >= ($${params.length}::date)::timestamp AT TIME ZONE 'Asia/Shanghai'`, + ) + } + if (createdTo) { + params.push(createdTo) + filters.push( + `wo.created_at < (($${params.length}::date + 1)::timestamp AT TIME ZONE 'Asia/Shanghai')`, + ) + } + 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(workerSharingId)} = 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 (acceptanceMode) { + params.push(normalizeWorkOrderAcceptanceMode(acceptanceMode)) + filters.push(`wo.acceptance_mode = $${params.length}`) + } + if (giftPhase) { + const normalizedGiftPhase = String(giftPhase).trim() + if (normalizedGiftPhase === WORK_ORDER_GIFT_PHASE.GIFT_READY) { + // gift_ready 是冷却到点后的实时阶段,由时间条件计算而不是落库状态。 + filters.push(`( + wo.acceptance_mode = 'friend_gift' + AND wo.gift_phase = 'friend_countdown' + AND wo.gift_available_at IS NOT NULL + AND wo.gift_available_at <= NOW() + )`) + } else if (Object.values(WORK_ORDER_GIFT_PHASE).includes(normalizedGiftPhase as never)) { + params.push(normalizedGiftPhase) + const giftPhaseFilter = `wo.acceptance_mode = 'friend_gift' AND wo.gift_phase = $${params.length}` + // 待打手提交资料仅代表已接单后的赠送流程;待抢单/未分配模板单不应混入该筛选。 + filters.push( + normalizedGiftPhase === WORK_ORDER_GIFT_PHASE.MATERIAL_REQUIRED + ? `(${giftPhaseFilter} AND wo.assigned_worker_id IS NOT NULL AND wo.status IN ('in_progress', 'problem'))` + : giftPhaseFilter, + ) + } + } + if (hallCandidateLimit > 0) { + params.push(hallCandidateLimit) + const limitPlaceholder = `$${params.length}` + params.push(JSON.stringify(hallCategoryLimits)) + filters.push(`wo.id IN (${buildHallCandidateIdsSql(limitPlaceholder, `$${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, categoryLimitsPlaceholder: string) { + return ` + SELECT ranked.id + FROM ( + SELECT candidate.id, candidate.category_id, candidate.hall_queued_at, + ROW_NUMBER() OVER ( + PARTITION BY candidate.category_id + ORDER BY candidate.hall_queued_at ASC NULLS LAST, candidate.id ASC + ) AS category_rank + 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 + ) + ) ranked + WHERE ranked.category_rank <= COALESCE( + NULLIF(${categoryLimitsPlaceholder}::jsonb ->> ranked.category_id::text, '')::int, + ${limitPlaceholder} + ) + ORDER BY ranked.hall_queued_at ASC NULLS LAST, ranked.id ASC + LIMIT ${limitPlaceholder} + ` +} + +/** 抢单操作在事务中复核容量与当前打手的可见延迟,避免通过订单 ID 绕过大厅队列。 */ +export async function isWorkOrderWithinHallCapacityWithClient( + client: PoolClient, + input: { + workOrderId: number + hallCandidateLimit: number + hallCategoryLimits?: Record + 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', '$3')}) + AND ($4::timestamptz IS NULL OR wo.hall_queued_at IS NULL OR wo.hall_queued_at <= $4) + ) AS visible + `, + [ + input.workOrderId, + input.hallCandidateLimit, + JSON.stringify(input.hallCategoryLimits || {}), + input.visibleAfterIso ? input.visibleAfterIso : null, + ], + ) + return result.rows[0]?.visible === true +} diff --git a/test-fix-artifacts/DIFF_FILE.patch b/test-fix-artifacts/DIFF_FILE.patch new file mode 100644 index 00000000..e78b528c --- /dev/null +++ b/test-fix-artifacts/DIFF_FILE.patch @@ -0,0 +1,9 @@ +Test fix diff summary + +1. worker-platform/work-order-query-repo.ts + - Trim keyword input inside buildWorkOrderWhere before adding LIKE filters. + - Whitespace-only keywords now add no filter and no SQL parameter. + +2. db/postgres.integration.test.ts + - Use a unique regular table for the pool-concurrency test. + - Drop the table in finally so concurrent pool connections can see it and cleanup is deterministic. diff --git a/test-fix-artifacts/MODIFIED_FILE.ts b/test-fix-artifacts/MODIFIED_FILE.ts new file mode 100644 index 00000000..b84151c2 --- /dev/null +++ b/test-fix-artifacts/MODIFIED_FILE.ts @@ -0,0 +1,961 @@ +import { query } from '../../db/client.js' +import type { ListInput, WorkOrderRow, WorkOrderStatistics } from './types.js' +import type { PoolClient } from 'pg' +import { + normalizeWorkOrderAcceptanceMode, + WORK_ORDER_GIFT_PHASE, +} from '../../domain/work-order-acceptance-mode.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, + cancel_info.cancel_reason, + cancel_info.cancel_event_type, + cancel_info.cancel_actor_type, + cancel_info.cancel_actor_id, + cancel_info.cancelled_at + 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 + ev.event_type AS cancel_event_type, + ev.actor_type AS cancel_actor_type, + ev.actor_id AS cancel_actor_id, + ev.created_at AS cancelled_at, + BTRIM(COALESCE( + NULLIF(BTRIM(COALESCE(ev.payload_json->>'reason', '')), ''), + NULLIF(BTRIM(COALESCE(ev.payload_json->>'note', '')), ''), + CASE ev.event_type + WHEN 'timeout_cancel_release' THEN '任务超时未完成,系统自动取消并退还押金' + WHEN 'timeout_cancel_deduct' THEN '任务超时未完成,系统自动取消并扣除押金' + ELSE '' + END + )) AS cancel_reason + FROM work_order_events ev + WHERE ev.work_order_id = wo.id + AND wo.status = 'cancelled' + AND ev.event_type IN ( + 'cancelled_by_admin', + 'problem_cancel_release', + 'problem_cancel_deduct', + 'timeout_cancel_release', + 'timeout_cancel_deduct' + ) + ORDER BY ev.id DESC + LIMIT 1 + ) cancel_info ON TRUE + 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 +` + +// Keep the page and total in one filtered query. This avoids running the same +// joins and share aggregation twice for every list request. +const LIST_WORK_ORDER_SELECT = WORK_ORDER_SELECT.replace( + ' SELECT\n', + ' SELECT\n COUNT(*) OVER()::int AS total_count,\n', +) + +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 +} + +type TimeoutCountCacheEntry = { + expiresAt: number + counts: Map +} + +// 后台打手列表会在短时间内重复请求同一页;超时次数只是展示字段,允许短 TTL。 +const TIMEOUT_COUNT_CACHE_TTL_MS = 5_000 +const TIMEOUT_COUNT_CACHE_MAX_ENTRIES = 256 +const timeoutCountCache = new Map() + +/** 后台按母工单拼单参与情况计算业务状态,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', '') + )` +} + +/** + * 打手订单个人状态:拼单按个人份额状态计算。 + * 传入 workerId 时,代练中但存在待审核取消申请或未处理反馈的订单归入问题单, + * 客服驳回/解决后自动回到代练中(计算型归类,无状态迁移)。 + */ +export function workerOrderStatusExpression(workerId = 0, workerShareAlias = 'worker_share') { + const baseExpression = `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` + const normalizedWorkerId = Math.max(0, Math.floor(Number(workerId) || 0)) + if (normalizedWorkerId <= 0) return baseExpression + const pendingHandlingExists = `( + EXISTS ( + SELECT 1 + FROM worker_order_cancel_requests wcr + WHERE wcr.work_order_id = wo.id + AND wcr.worker_id = ${normalizedWorkerId} + AND wcr.status = 'pending' + ) OR EXISTS ( + SELECT 1 + FROM worker_order_feedbacks wf + WHERE wf.work_order_id = wo.id + AND wf.worker_id = ${normalizedWorkerId} + AND wf.status = 'pending' + ) + )` + return `CASE + WHEN (${baseExpression}) = 'in_progress' AND ${pendingHandlingExists} THEN 'problem' + ELSE ${baseExpression} + 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, + hallCategoryLimits = {}, + visibleAfterIso = '', + excludeFilledSharing = false, + vipEvidencePending = false, + acceptanceMode = '', + giftPhase = '', + useOperationalStatus = false, + pinnedFirst = false, + sort = 'id_desc', + createdFrom = '', + createdTo = '', +}: 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, + hallCategoryLimits, + visibleAfterIso, + excludeFilledSharing, + vipEvidencePending, + acceptanceMode, + giftPhase, + useOperationalStatus, + createdFrom, + createdTo, + }) + 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( + `${LIST_WORK_ORDER_SELECT} + ${joins} + ${whereClause} + ORDER BY ${effectiveOrderBy} + LIMIT $${params.length - 1} OFFSET $${params.length}`, + params, + ) + // A page beyond the end has no row from which to read the window count. + // Keep that uncommon pagination edge case correct without restoring a + // second query for normal requests. + let total = Number( + (itemsResult.rows[0] as WorkOrderRow & { total_count?: number })?.total_count || 0, + ) + if (itemsResult.rows.length === 0 && page > 1) { + const totalResult = await query<{ total: number }>( + `SELECT COUNT(*)::int AS total FROM work_orders wo ${joins} ${whereClause}`, + params.slice(0, -2), + ) + total = Number(totalResult.rows[0]?.total || 0) + } + return { + items: itemsResult.rows, + total, + } +} + +/** 按筛选条件汇总工单,统计范围不受当前分页影响。 */ +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 visible_orders AS ( + SELECT wo.id AS work_order_id + FROM work_orders wo + WHERE wo.assigned_worker_id = $1 + UNION + SELECT wo.id AS work_order_id + FROM work_orders wo + WHERE 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' + ) + ) + UNION + SELECT wos.work_order_id + FROM work_order_shares wos + WHERE wos.worker_id = $1 + AND wos.status != 'cancelled' + ), + worker_base AS ( + SELECT + wo.id, + wo.status AS wo_status, + wos.status AS wos_status, + 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, + wo.acceptance_json + FROM work_orders wo + INNER JOIN visible_orders vo ON vo.work_order_id = wo.id + LEFT JOIN work_order_shares wos + ON wos.work_order_id = wo.id + AND wos.worker_id = $1 + AND wos.status != 'cancelled' + ), + worker_orders AS ( + SELECT + CASE + WHEN ( + CASE + WHEN is_worker_share 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 + ) = 'in_progress' AND ( + EXISTS ( + SELECT 1 + FROM worker_order_cancel_requests wcr + WHERE wcr.work_order_id = worker_base.id + AND wcr.worker_id = $1 + AND wcr.status = 'pending' + ) OR EXISTS ( + SELECT 1 + FROM worker_order_feedbacks wf + WHERE wf.work_order_id = worker_base.id + AND wf.worker_id = $1 + AND wf.status = 'pending' + ) + ) THEN 'problem' + ELSE CASE + WHEN is_worker_share 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 + END AS status, + acceptance_json, + is_worker_share, + worker_settlement_amount + FROM worker_base + ) + 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 listOverdueWorkOrders({ + limit = 50, + workerId = 0, +}: { limit?: number; workerId?: number } = {}): Promise { + const params: unknown[] = [] + let workerClause = '' + if (workerId > 0) { + params.push(workerId) + workerClause = `AND wo.assigned_worker_id = $${params.length}` + } + params.push(limit) + const result = await query( + `${WORK_ORDER_SELECT} + WHERE wo.status = 'in_progress' + AND wo.deadline_at IS NOT NULL + AND wo.deadline_at < NOW() + AND NOT (wo.acceptance_mode = 'friend_gift' AND wo.gift_phase = 'material_submitted') + ${workerClause} + ORDER BY wo.deadline_at ASC + LIMIT $${params.length}`, + params, + ) + return result.rows +} + +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 countWorkerActiveOrdersWithClient( + client: PoolClient, + workerId: number, +): Promise { + const result = await client.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 + `, + [workerId], + ) + return Number(result.rows[0]?.total || 0) +} + +export async function countTimeoutEventsByWorkerIds( + workerIds: number[], +): Promise> { + const uniqueIds = [...new Set(workerIds.map((id) => Number(id)).filter((id) => id > 0))] + if (uniqueIds.length === 0) return new Map() + uniqueIds.sort((left, right) => left - right) + const cacheKey = uniqueIds.join(',') + const now = Date.now() + const cached = timeoutCountCache.get(cacheKey) + if (cached && cached.expiresAt > now) { + return new Map(cached.counts) + } + + const counts = new Map() + 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)) + } + } + if (timeoutCountCache.size >= TIMEOUT_COUNT_CACHE_MAX_ENTRIES) { + const oldestKey = timeoutCountCache.keys().next().value + if (oldestKey) timeoutCountCache.delete(oldestKey) + } + timeoutCountCache.set(cacheKey, { + expiresAt: now + TIMEOUT_COUNT_CACHE_TTL_MS, + counts: new Map(counts), + }) + 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 countWorkerCancellationsSince( + workerId: number | string, + sinceIso: string, +): Promise { + const result = await query<{ total: number }>( + ` + SELECT COUNT(*)::int AS total + FROM work_order_events + WHERE actor_type = 'worker' + AND actor_id = $1 + AND event_type = 'cancelled_by_worker' + AND created_at >= $2 + `, + [String(workerId), sinceIso], + ) + return Number(result.rows[0]?.total || 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, + hallCategoryLimits = {}, + visibleAfterIso = '', + excludeFilledSharing = false, + vipEvidencePending = false, + acceptanceMode = '', + giftPhase = '', + useOperationalStatus = false, + createdFrom = '', + createdTo = '', +}: ListInput) { + const filters: string[] = [] + const params: unknown[] = [] + const joins: string[] = [] + const normalizedKeyword = String(keyword || '').trim() + if (createdFrom) { + params.push(createdFrom) + filters.push( + `wo.created_at >= ($${params.length}::date)::timestamp AT TIME ZONE 'Asia/Shanghai'`, + ) + } + if (createdTo) { + params.push(createdTo) + filters.push( + `wo.created_at < (($${params.length}::date + 1)::timestamp AT TIME ZONE 'Asia/Shanghai')`, + ) + } + 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(workerSharingId)} = 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 (normalizedKeyword) { + params.push(`%${normalizedKeyword}%`) + 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 (acceptanceMode) { + params.push(normalizeWorkOrderAcceptanceMode(acceptanceMode)) + filters.push(`wo.acceptance_mode = $${params.length}`) + } + if (giftPhase) { + const normalizedGiftPhase = String(giftPhase).trim() + if (normalizedGiftPhase === WORK_ORDER_GIFT_PHASE.GIFT_READY) { + // gift_ready 是冷却到点后的实时阶段,由时间条件计算而不是落库状态。 + filters.push(`( + wo.acceptance_mode = 'friend_gift' + AND wo.gift_phase = 'friend_countdown' + AND wo.gift_available_at IS NOT NULL + AND wo.gift_available_at <= NOW() + )`) + } else if (Object.values(WORK_ORDER_GIFT_PHASE).includes(normalizedGiftPhase as never)) { + params.push(normalizedGiftPhase) + const giftPhaseFilter = `wo.acceptance_mode = 'friend_gift' AND wo.gift_phase = $${params.length}` + // 待打手提交资料仅代表已接单后的赠送流程;待抢单/未分配模板单不应混入该筛选。 + filters.push( + normalizedGiftPhase === WORK_ORDER_GIFT_PHASE.MATERIAL_REQUIRED + ? `(${giftPhaseFilter} AND wo.assigned_worker_id IS NOT NULL AND wo.status IN ('in_progress', 'problem'))` + : giftPhaseFilter, + ) + } + } + if (hallCandidateLimit > 0) { + params.push(hallCandidateLimit) + const limitPlaceholder = `$${params.length}` + params.push(JSON.stringify(hallCategoryLimits)) + filters.push(`wo.id IN (${buildHallCandidateIdsSql(limitPlaceholder, `$${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, categoryLimitsPlaceholder: string) { + return ` + SELECT ranked.id + FROM ( + SELECT candidate.id, candidate.category_id, candidate.hall_queued_at, + ROW_NUMBER() OVER ( + PARTITION BY candidate.category_id + ORDER BY candidate.hall_queued_at ASC NULLS LAST, candidate.id ASC + ) AS category_rank + 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 + ) + ) ranked + WHERE ranked.category_rank <= COALESCE( + NULLIF(${categoryLimitsPlaceholder}::jsonb ->> ranked.category_id::text, '')::int, + ${limitPlaceholder} + ) + ORDER BY ranked.hall_queued_at ASC NULLS LAST, ranked.id ASC + LIMIT ${limitPlaceholder} + ` +} + +/** 抢单操作在事务中复核容量与当前打手的可见延迟,避免通过订单 ID 绕过大厅队列。 */ +export async function isWorkOrderWithinHallCapacityWithClient( + client: PoolClient, + input: { + workOrderId: number + hallCandidateLimit: number + hallCategoryLimits?: Record + 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', '$3')}) + AND ($4::timestamptz IS NULL OR wo.hall_queued_at IS NULL OR wo.hall_queued_at <= $4) + ) AS visible + `, + [ + input.workOrderId, + input.hallCandidateLimit, + JSON.stringify(input.hallCategoryLimits || {}), + input.visibleAfterIso ? input.visibleAfterIso : null, + ], + ) + return result.rows[0]?.visible === true +} diff --git a/test-fix-artifacts/ROLLBACK.sh b/test-fix-artifacts/ROLLBACK.sh new file mode 100755 index 00000000..ef84396f --- /dev/null +++ b/test-fix-artifacts/ROLLBACK.sh @@ -0,0 +1,9 @@ +#!/usr/bin/env bash +set -euo pipefail + +ROOT="${1:?usage: ROLLBACK.sh }" +SCRIPT_DIR="$(CDPATH= cd -- "$(dirname -- "$0")" && pwd)" + +cp "$SCRIPT_DIR/BASELINE_worker_query.ts" "$ROOT/apps/backend/src/repositories/worker-platform/work-order-query-repo.ts" +cp "$SCRIPT_DIR/BASELINE_postgres_integration.test.ts" "$ROOT/apps/backend/src/db/postgres.integration.test.ts" +printf 'restored worker query and PostgreSQL integration test baselines in %s\n' "$ROOT" diff --git a/test-fix-artifacts/VERIFICATION.txt b/test-fix-artifacts/VERIFICATION.txt new file mode 100644 index 00000000..c5aa2ab8 --- /dev/null +++ b/test-fix-artifacts/VERIFICATION.txt @@ -0,0 +1,27 @@ +Test failure fix verification + +Changed branch/fields: +- buildWorkOrderWhere keyword normalization: whitespace-only input no longer emits LIKE SQL. +- PostgreSQL concurrency test table scope: regular unique table visible to all pool connections, dropped in finally. + +Artifacts: +- MODIFIED_FILE: /Users/yml/codes/order_site/test-fix-artifacts/MODIFIED_FILE.ts +- DIFF_FILE: /Users/yml/codes/order_site/test-fix-artifacts/DIFF_FILE.patch +- VERIFICATION.txt: /Users/yml/codes/order_site/test-fix-artifacts/VERIFICATION.txt +- ROLLBACK.sh: /Users/yml/codes/order_site/test-fix-artifacts/ROLLBACK.sh + +BASELINE +Command: git show HEAD:apps/backend/src/repositories/worker-platform/work-order-query-repo.ts +Input: HEAD source before this fix +Literal result: SHA-256 ff715033b3939b18788dab93e13fa0b3eaf054e469eb721c1b8fa7611554b13e; exit status 0 + +MODIFIED +Command: TEST_POSTGRES_URL=postgres://postgres:postgres@127.0.0.1:5432/order_site npm run check +Input: modified worker query and PostgreSQL integration test sources +Literal result: format passed; SQL guard passed; lint passed; typecheck passed; tests 377; pass 377; fail 0; skipped 0; build completed; exit status 0 + +ROLLBACK +Command: test-fix-artifacts/ROLLBACK.sh test-fix-artifacts/rollback-copy +Input: independent copy containing MODIFIED_FILE.ts and modified PostgreSQL integration test +Literal result: restored worker query and PostgreSQL integration test baselines; restored worker query SHA-256 ff715033b3939b18788dab93e13fa0b3eaf054e469eb721c1b8fa7611554b13e; exit status 0 +Restored behavior/status: rollback-copy worker query matches baseline; working source remains fixed. diff --git a/test-fix-artifacts/rollback-copy/apps/backend/src/db/postgres.integration.test.ts b/test-fix-artifacts/rollback-copy/apps/backend/src/db/postgres.integration.test.ts new file mode 100644 index 00000000..c9d28eeb --- /dev/null +++ b/test-fix-artifacts/rollback-copy/apps/backend/src/db/postgres.integration.test.ts @@ -0,0 +1,60 @@ +import assert from 'node:assert/strict' +import test from 'node:test' +import { Pool } from 'pg' + +const connectionString = process.env.TEST_POSTGRES_URL + +test('real PostgreSQL transaction rolls back atomically', { skip: !connectionString }, async () => { + const pool = new Pool({ connectionString }) + const client = await pool.connect() + const tableName = `integration_tx_${Date.now()}` + try { + await client.query(`CREATE TEMP TABLE ${tableName} (id integer primary key, value text)`) + await client.query('BEGIN') + await client.query(`INSERT INTO ${tableName} (id, value) VALUES (1, 'before-error')`) + try { + await client.query(`INSERT INTO ${tableName} (id, value) VALUES (1, 'duplicate')`) + } catch { + await client.query('ROLLBACK') + } + const result = await client.query(`SELECT COUNT(*)::int AS count FROM ${tableName}`) + assert.equal(result.rows[0].count, 0) + } finally { + client.release() + await pool.end() + } +}) + +test( + 'real PostgreSQL enforces callback replay idempotency under concurrency', + { + skip: !connectionString, + }, + async () => { + const pool = new Pool({ connectionString }) + const tableName = `integration_webhook_${Date.now()}` + try { + await pool.query( + `CREATE TEMP TABLE ${tableName} ( + id serial primary key, + provider text not null, + event_id text not null, + unique(provider, event_id) + )`, + ) + await Promise.all( + Array.from({ length: 8 }, () => + pool.query( + `INSERT INTO ${tableName} (provider, event_id) + VALUES ('integration', 'replay-1') + ON CONFLICT (provider, event_id) DO NOTHING`, + ), + ), + ) + const result = await pool.query(`SELECT COUNT(*)::int AS count FROM ${tableName}`) + assert.equal(result.rows[0].count, 1) + } finally { + await pool.end() + } + }, +) diff --git a/test-fix-artifacts/rollback-copy/apps/backend/src/repositories/worker-platform/work-order-query-repo.ts b/test-fix-artifacts/rollback-copy/apps/backend/src/repositories/worker-platform/work-order-query-repo.ts new file mode 100644 index 00000000..15a19d68 --- /dev/null +++ b/test-fix-artifacts/rollback-copy/apps/backend/src/repositories/worker-platform/work-order-query-repo.ts @@ -0,0 +1,960 @@ +import { query } from '../../db/client.js' +import type { ListInput, WorkOrderRow, WorkOrderStatistics } from './types.js' +import type { PoolClient } from 'pg' +import { + normalizeWorkOrderAcceptanceMode, + WORK_ORDER_GIFT_PHASE, +} from '../../domain/work-order-acceptance-mode.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, + cancel_info.cancel_reason, + cancel_info.cancel_event_type, + cancel_info.cancel_actor_type, + cancel_info.cancel_actor_id, + cancel_info.cancelled_at + 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 + ev.event_type AS cancel_event_type, + ev.actor_type AS cancel_actor_type, + ev.actor_id AS cancel_actor_id, + ev.created_at AS cancelled_at, + BTRIM(COALESCE( + NULLIF(BTRIM(COALESCE(ev.payload_json->>'reason', '')), ''), + NULLIF(BTRIM(COALESCE(ev.payload_json->>'note', '')), ''), + CASE ev.event_type + WHEN 'timeout_cancel_release' THEN '任务超时未完成,系统自动取消并退还押金' + WHEN 'timeout_cancel_deduct' THEN '任务超时未完成,系统自动取消并扣除押金' + ELSE '' + END + )) AS cancel_reason + FROM work_order_events ev + WHERE ev.work_order_id = wo.id + AND wo.status = 'cancelled' + AND ev.event_type IN ( + 'cancelled_by_admin', + 'problem_cancel_release', + 'problem_cancel_deduct', + 'timeout_cancel_release', + 'timeout_cancel_deduct' + ) + ORDER BY ev.id DESC + LIMIT 1 + ) cancel_info ON TRUE + 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 +` + +// Keep the page and total in one filtered query. This avoids running the same +// joins and share aggregation twice for every list request. +const LIST_WORK_ORDER_SELECT = WORK_ORDER_SELECT.replace( + ' SELECT\n', + ' SELECT\n COUNT(*) OVER()::int AS total_count,\n', +) + +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 +} + +type TimeoutCountCacheEntry = { + expiresAt: number + counts: Map +} + +// 后台打手列表会在短时间内重复请求同一页;超时次数只是展示字段,允许短 TTL。 +const TIMEOUT_COUNT_CACHE_TTL_MS = 5_000 +const TIMEOUT_COUNT_CACHE_MAX_ENTRIES = 256 +const timeoutCountCache = new Map() + +/** 后台按母工单拼单参与情况计算业务状态,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', '') + )` +} + +/** + * 打手订单个人状态:拼单按个人份额状态计算。 + * 传入 workerId 时,代练中但存在待审核取消申请或未处理反馈的订单归入问题单, + * 客服驳回/解决后自动回到代练中(计算型归类,无状态迁移)。 + */ +export function workerOrderStatusExpression(workerId = 0, workerShareAlias = 'worker_share') { + const baseExpression = `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` + const normalizedWorkerId = Math.max(0, Math.floor(Number(workerId) || 0)) + if (normalizedWorkerId <= 0) return baseExpression + const pendingHandlingExists = `( + EXISTS ( + SELECT 1 + FROM worker_order_cancel_requests wcr + WHERE wcr.work_order_id = wo.id + AND wcr.worker_id = ${normalizedWorkerId} + AND wcr.status = 'pending' + ) OR EXISTS ( + SELECT 1 + FROM worker_order_feedbacks wf + WHERE wf.work_order_id = wo.id + AND wf.worker_id = ${normalizedWorkerId} + AND wf.status = 'pending' + ) + )` + return `CASE + WHEN (${baseExpression}) = 'in_progress' AND ${pendingHandlingExists} THEN 'problem' + ELSE ${baseExpression} + 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, + hallCategoryLimits = {}, + visibleAfterIso = '', + excludeFilledSharing = false, + vipEvidencePending = false, + acceptanceMode = '', + giftPhase = '', + useOperationalStatus = false, + pinnedFirst = false, + sort = 'id_desc', + createdFrom = '', + createdTo = '', +}: 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, + hallCategoryLimits, + visibleAfterIso, + excludeFilledSharing, + vipEvidencePending, + acceptanceMode, + giftPhase, + useOperationalStatus, + createdFrom, + createdTo, + }) + 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( + `${LIST_WORK_ORDER_SELECT} + ${joins} + ${whereClause} + ORDER BY ${effectiveOrderBy} + LIMIT $${params.length - 1} OFFSET $${params.length}`, + params, + ) + // A page beyond the end has no row from which to read the window count. + // Keep that uncommon pagination edge case correct without restoring a + // second query for normal requests. + let total = Number( + (itemsResult.rows[0] as WorkOrderRow & { total_count?: number })?.total_count || 0, + ) + if (itemsResult.rows.length === 0 && page > 1) { + const totalResult = await query<{ total: number }>( + `SELECT COUNT(*)::int AS total FROM work_orders wo ${joins} ${whereClause}`, + params.slice(0, -2), + ) + total = Number(totalResult.rows[0]?.total || 0) + } + return { + items: itemsResult.rows, + total, + } +} + +/** 按筛选条件汇总工单,统计范围不受当前分页影响。 */ +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 visible_orders AS ( + SELECT wo.id AS work_order_id + FROM work_orders wo + WHERE wo.assigned_worker_id = $1 + UNION + SELECT wo.id AS work_order_id + FROM work_orders wo + WHERE 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' + ) + ) + UNION + SELECT wos.work_order_id + FROM work_order_shares wos + WHERE wos.worker_id = $1 + AND wos.status != 'cancelled' + ), + worker_base AS ( + SELECT + wo.id, + wo.status AS wo_status, + wos.status AS wos_status, + 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, + wo.acceptance_json + FROM work_orders wo + INNER JOIN visible_orders vo ON vo.work_order_id = wo.id + LEFT JOIN work_order_shares wos + ON wos.work_order_id = wo.id + AND wos.worker_id = $1 + AND wos.status != 'cancelled' + ), + worker_orders AS ( + SELECT + CASE + WHEN ( + CASE + WHEN is_worker_share 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 + ) = 'in_progress' AND ( + EXISTS ( + SELECT 1 + FROM worker_order_cancel_requests wcr + WHERE wcr.work_order_id = worker_base.id + AND wcr.worker_id = $1 + AND wcr.status = 'pending' + ) OR EXISTS ( + SELECT 1 + FROM worker_order_feedbacks wf + WHERE wf.work_order_id = worker_base.id + AND wf.worker_id = $1 + AND wf.status = 'pending' + ) + ) THEN 'problem' + ELSE CASE + WHEN is_worker_share 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 + END AS status, + acceptance_json, + is_worker_share, + worker_settlement_amount + FROM worker_base + ) + 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 listOverdueWorkOrders({ + limit = 50, + workerId = 0, +}: { limit?: number; workerId?: number } = {}): Promise { + const params: unknown[] = [] + let workerClause = '' + if (workerId > 0) { + params.push(workerId) + workerClause = `AND wo.assigned_worker_id = $${params.length}` + } + params.push(limit) + const result = await query( + `${WORK_ORDER_SELECT} + WHERE wo.status = 'in_progress' + AND wo.deadline_at IS NOT NULL + AND wo.deadline_at < NOW() + AND NOT (wo.acceptance_mode = 'friend_gift' AND wo.gift_phase = 'material_submitted') + ${workerClause} + ORDER BY wo.deadline_at ASC + LIMIT $${params.length}`, + params, + ) + return result.rows +} + +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 countWorkerActiveOrdersWithClient( + client: PoolClient, + workerId: number, +): Promise { + const result = await client.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 + `, + [workerId], + ) + return Number(result.rows[0]?.total || 0) +} + +export async function countTimeoutEventsByWorkerIds( + workerIds: number[], +): Promise> { + const uniqueIds = [...new Set(workerIds.map((id) => Number(id)).filter((id) => id > 0))] + if (uniqueIds.length === 0) return new Map() + uniqueIds.sort((left, right) => left - right) + const cacheKey = uniqueIds.join(',') + const now = Date.now() + const cached = timeoutCountCache.get(cacheKey) + if (cached && cached.expiresAt > now) { + return new Map(cached.counts) + } + + const counts = new Map() + 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)) + } + } + if (timeoutCountCache.size >= TIMEOUT_COUNT_CACHE_MAX_ENTRIES) { + const oldestKey = timeoutCountCache.keys().next().value + if (oldestKey) timeoutCountCache.delete(oldestKey) + } + timeoutCountCache.set(cacheKey, { + expiresAt: now + TIMEOUT_COUNT_CACHE_TTL_MS, + counts: new Map(counts), + }) + 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 countWorkerCancellationsSince( + workerId: number | string, + sinceIso: string, +): Promise { + const result = await query<{ total: number }>( + ` + SELECT COUNT(*)::int AS total + FROM work_order_events + WHERE actor_type = 'worker' + AND actor_id = $1 + AND event_type = 'cancelled_by_worker' + AND created_at >= $2 + `, + [String(workerId), sinceIso], + ) + return Number(result.rows[0]?.total || 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, + hallCategoryLimits = {}, + visibleAfterIso = '', + excludeFilledSharing = false, + vipEvidencePending = false, + acceptanceMode = '', + giftPhase = '', + useOperationalStatus = false, + createdFrom = '', + createdTo = '', +}: ListInput) { + const filters: string[] = [] + const params: unknown[] = [] + const joins: string[] = [] + if (createdFrom) { + params.push(createdFrom) + filters.push( + `wo.created_at >= ($${params.length}::date)::timestamp AT TIME ZONE 'Asia/Shanghai'`, + ) + } + if (createdTo) { + params.push(createdTo) + filters.push( + `wo.created_at < (($${params.length}::date + 1)::timestamp AT TIME ZONE 'Asia/Shanghai')`, + ) + } + 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(workerSharingId)} = 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 (acceptanceMode) { + params.push(normalizeWorkOrderAcceptanceMode(acceptanceMode)) + filters.push(`wo.acceptance_mode = $${params.length}`) + } + if (giftPhase) { + const normalizedGiftPhase = String(giftPhase).trim() + if (normalizedGiftPhase === WORK_ORDER_GIFT_PHASE.GIFT_READY) { + // gift_ready 是冷却到点后的实时阶段,由时间条件计算而不是落库状态。 + filters.push(`( + wo.acceptance_mode = 'friend_gift' + AND wo.gift_phase = 'friend_countdown' + AND wo.gift_available_at IS NOT NULL + AND wo.gift_available_at <= NOW() + )`) + } else if (Object.values(WORK_ORDER_GIFT_PHASE).includes(normalizedGiftPhase as never)) { + params.push(normalizedGiftPhase) + const giftPhaseFilter = `wo.acceptance_mode = 'friend_gift' AND wo.gift_phase = $${params.length}` + // 待打手提交资料仅代表已接单后的赠送流程;待抢单/未分配模板单不应混入该筛选。 + filters.push( + normalizedGiftPhase === WORK_ORDER_GIFT_PHASE.MATERIAL_REQUIRED + ? `(${giftPhaseFilter} AND wo.assigned_worker_id IS NOT NULL AND wo.status IN ('in_progress', 'problem'))` + : giftPhaseFilter, + ) + } + } + if (hallCandidateLimit > 0) { + params.push(hallCandidateLimit) + const limitPlaceholder = `$${params.length}` + params.push(JSON.stringify(hallCategoryLimits)) + filters.push(`wo.id IN (${buildHallCandidateIdsSql(limitPlaceholder, `$${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, categoryLimitsPlaceholder: string) { + return ` + SELECT ranked.id + FROM ( + SELECT candidate.id, candidate.category_id, candidate.hall_queued_at, + ROW_NUMBER() OVER ( + PARTITION BY candidate.category_id + ORDER BY candidate.hall_queued_at ASC NULLS LAST, candidate.id ASC + ) AS category_rank + 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 + ) + ) ranked + WHERE ranked.category_rank <= COALESCE( + NULLIF(${categoryLimitsPlaceholder}::jsonb ->> ranked.category_id::text, '')::int, + ${limitPlaceholder} + ) + ORDER BY ranked.hall_queued_at ASC NULLS LAST, ranked.id ASC + LIMIT ${limitPlaceholder} + ` +} + +/** 抢单操作在事务中复核容量与当前打手的可见延迟,避免通过订单 ID 绕过大厅队列。 */ +export async function isWorkOrderWithinHallCapacityWithClient( + client: PoolClient, + input: { + workOrderId: number + hallCandidateLimit: number + hallCategoryLimits?: Record + 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', '$3')}) + AND ($4::timestamptz IS NULL OR wo.hall_queued_at IS NULL OR wo.hall_queued_at <= $4) + ) AS visible + `, + [ + input.workOrderId, + input.hallCandidateLimit, + JSON.stringify(input.hallCategoryLimits || {}), + input.visibleAfterIso ? input.visibleAfterIso : null, + ], + ) + return result.rows[0]?.visible === true +}