修复:恢复 PostgreSQL 集成测试并规范工单搜索条件

This commit is contained in:
yml2213
2026-08-30 17:04:35 +08:00
parent 25b89bb503
commit 535cefd4de
10 changed files with 3051 additions and 3 deletions
@@ -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()
}
},
@@ -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) {
@@ -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()
}
},
)
+960
View File
@@ -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<number, number>
}
// 后台打手列表会在短时间内重复请求同一页;超时次数只是展示字段,允许短 TTL。
const TIMEOUT_COUNT_CACHE_TTL_MS = 5_000
const TIMEOUT_COUNT_CACHE_MAX_ENTRIES = 256
const timeoutCountCache = new Map<string, TimeoutCountCacheEntry>()
/** 后台按母工单拼单参与情况计算业务状态,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<WorkOrderRow | null> {
const result = await query<WorkOrderRow>(
`${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<WorkOrderRow | null> {
const result = await query<WorkOrderRow>(`${WORK_ORDER_SELECT} WHERE wo.id = $1 LIMIT 1`, [
Number(workOrderId),
])
return result.rows[0] || null
}
export async function findPendingMaterialWorkOrderByPlatformOrderId(
platformOrderId: string,
): Promise<WorkOrderRow | null> {
const result = await query<WorkOrderRow>(
`${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<WorkOrderRow[]> {
const result = await query<WorkOrderRow>(
`${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<WorkOrderRow[]> {
const result = await query<WorkOrderRow>(
`${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<WorkOrderRow>(
`${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<WorkOrderStatistics> {
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<WorkOrderStatisticsRow>(
`
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<WorkerOrderOverviewRow>(
`
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<WorkOrderRow[]> {
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<WorkOrderRow>(
`${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<number> {
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<number> {
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<Map<number, number>> {
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<number, number>()
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<number> {
const counts = await countTimeoutEventsByWorkerIds([Number(workerId)])
return counts.get(Number(workerId)) || 0
}
export async function countWorkerCancellationsSince(
workerId: number | string,
sinceIso: string,
): Promise<number> {
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<WorkOrderRow | null> {
const result = await client.query<WorkOrderRow>(`${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<string, number>
visibleAfterIso?: string | undefined
},
) {
const result = await client.query<{ visible: boolean }>(
`
SELECT EXISTS(
SELECT 1
FROM work_orders wo
WHERE wo.id = $1
AND wo.id IN (${buildHallCandidateIdsSql('$2', '$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
}
+9
View File
@@ -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.
+961
View File
@@ -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<number, number>
}
// 后台打手列表会在短时间内重复请求同一页;超时次数只是展示字段,允许短 TTL。
const TIMEOUT_COUNT_CACHE_TTL_MS = 5_000
const TIMEOUT_COUNT_CACHE_MAX_ENTRIES = 256
const timeoutCountCache = new Map<string, TimeoutCountCacheEntry>()
/** 后台按母工单拼单参与情况计算业务状态,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<WorkOrderRow | null> {
const result = await query<WorkOrderRow>(
`${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<WorkOrderRow | null> {
const result = await query<WorkOrderRow>(`${WORK_ORDER_SELECT} WHERE wo.id = $1 LIMIT 1`, [
Number(workOrderId),
])
return result.rows[0] || null
}
export async function findPendingMaterialWorkOrderByPlatformOrderId(
platformOrderId: string,
): Promise<WorkOrderRow | null> {
const result = await query<WorkOrderRow>(
`${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<WorkOrderRow[]> {
const result = await query<WorkOrderRow>(
`${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<WorkOrderRow[]> {
const result = await query<WorkOrderRow>(
`${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<WorkOrderRow>(
`${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<WorkOrderStatistics> {
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<WorkOrderStatisticsRow>(
`
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<WorkerOrderOverviewRow>(
`
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<WorkOrderRow[]> {
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<WorkOrderRow>(
`${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<number> {
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<number> {
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<Map<number, number>> {
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<number, number>()
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<number> {
const counts = await countTimeoutEventsByWorkerIds([Number(workerId)])
return counts.get(Number(workerId)) || 0
}
export async function countWorkerCancellationsSince(
workerId: number | string,
sinceIso: string,
): Promise<number> {
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<WorkOrderRow | null> {
const result = await client.query<WorkOrderRow>(`${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<string, number>
visibleAfterIso?: string | undefined
},
) {
const result = await client.query<{ visible: boolean }>(
`
SELECT EXISTS(
SELECT 1
FROM work_orders wo
WHERE wo.id = $1
AND wo.id IN (${buildHallCandidateIdsSql('$2', '$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
}
+9
View File
@@ -0,0 +1,9 @@
#!/usr/bin/env bash
set -euo pipefail
ROOT="${1:?usage: ROLLBACK.sh <workspace-copy>}"
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"
+27
View File
@@ -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.
@@ -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()
}
},
)
@@ -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<number, number>
}
// 后台打手列表会在短时间内重复请求同一页;超时次数只是展示字段,允许短 TTL。
const TIMEOUT_COUNT_CACHE_TTL_MS = 5_000
const TIMEOUT_COUNT_CACHE_MAX_ENTRIES = 256
const timeoutCountCache = new Map<string, TimeoutCountCacheEntry>()
/** 后台按母工单拼单参与情况计算业务状态,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<WorkOrderRow | null> {
const result = await query<WorkOrderRow>(
`${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<WorkOrderRow | null> {
const result = await query<WorkOrderRow>(`${WORK_ORDER_SELECT} WHERE wo.id = $1 LIMIT 1`, [
Number(workOrderId),
])
return result.rows[0] || null
}
export async function findPendingMaterialWorkOrderByPlatformOrderId(
platformOrderId: string,
): Promise<WorkOrderRow | null> {
const result = await query<WorkOrderRow>(
`${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<WorkOrderRow[]> {
const result = await query<WorkOrderRow>(
`${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<WorkOrderRow[]> {
const result = await query<WorkOrderRow>(
`${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<WorkOrderRow>(
`${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<WorkOrderStatistics> {
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<WorkOrderStatisticsRow>(
`
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<WorkerOrderOverviewRow>(
`
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<WorkOrderRow[]> {
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<WorkOrderRow>(
`${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<number> {
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<number> {
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<Map<number, number>> {
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<number, number>()
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<number> {
const counts = await countTimeoutEventsByWorkerIds([Number(workerId)])
return counts.get(Number(workerId)) || 0
}
export async function countWorkerCancellationsSince(
workerId: number | string,
sinceIso: string,
): Promise<number> {
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<WorkOrderRow | null> {
const result = await client.query<WorkOrderRow>(`${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<string, number>
visibleAfterIso?: string | undefined
},
) {
const result = await client.query<{ visible: boolean }>(
`
SELECT EXISTS(
SELECT 1
FROM work_orders wo
WHERE wo.id = $1
AND wo.id IN (${buildHallCandidateIdsSql('$2', '$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
}