拆分打手平台售后与订单模块

This commit is contained in:
yml2213
2026-08-21 15:17:28 +08:00
parent f4d32a7c07
commit 798ea583f8
24 changed files with 4473 additions and 4213 deletions
@@ -2,6 +2,13 @@ export * from './types.js'
export * from './shared.js'
export * from './worker-repo.js'
export * from './work-order-repo.js'
export * from './work-order-event-repo.js'
export * from './work-order-share-repo.js'
export * from './work-category-repo.js'
export * from './work-product-rule-repo.js'
export * from './worker-order-note-repo.js'
export * from './worker-order-feedback-repo.js'
export * from './worker-cancel-request-repo.js'
export * from './after-sales-repo.js'
export * from './product-match-repo.js'
export * from './sms-code-repo.js'
@@ -41,3 +41,22 @@ export function toPositiveInteger(value: unknown, fallback: number): number {
if (!Number.isFinite(parsed)) return fallback
return Math.max(0, Math.round(parsed))
}
export const WORK_ORDER_SHARE_SELECT = `
SELECT
wos.*,
wu.username AS worker_username,
wu.display_name AS worker_display_name
FROM work_order_shares wos
LEFT JOIN worker_users wu ON wu.id = wos.worker_id
`
/** 校验单个拼单份额能否撤单,并返回母工单应回退到的状态。 */
export function resolveWorkOrderShareCancellationStatus(
workOrderStatus: string,
shareStatus: string,
): 'open' | null {
if (!['open', 'pending_acceptance'].includes(workOrderStatus)) return null
if (!['joined', 'submitted'].includes(shareStatus)) return null
return 'open'
}
@@ -0,0 +1,84 @@
import { query } from '../../db/client.js'
import type { WorkCategoryRow } from './types.js'
export async function getWorkCategoryByKey(categoryKey: string): Promise<WorkCategoryRow | null> {
const result = await query<WorkCategoryRow>(
'SELECT * FROM work_categories WHERE category_key = $1 LIMIT 1',
[String(categoryKey || '').trim()],
)
return result.rows[0] || null
}
export async function getWorkCategoryById(
categoryId: number | string,
): Promise<WorkCategoryRow | null> {
const result = await query<WorkCategoryRow>(
'SELECT * FROM work_categories WHERE id = $1 LIMIT 1',
[Number(categoryId)],
)
return result.rows[0] || null
}
export async function listWorkCategories(): Promise<WorkCategoryRow[]> {
const result = await query<WorkCategoryRow>(
`SELECT * FROM work_categories WHERE status = 'active' ORDER BY sort_order ASC, id ASC`,
)
return result.rows
}
export async function listAllWorkCategories(): Promise<WorkCategoryRow[]> {
const result = await query<WorkCategoryRow>(
`SELECT * FROM work_categories ORDER BY sort_order ASC, id ASC`,
)
return result.rows
}
export async function upsertWorkCategory(input: {
categoryKey: string
name: string
sortOrder: number
status: string
now: string
}): Promise<WorkCategoryRow | null> {
const result = await query<WorkCategoryRow>(
`
INSERT INTO work_categories (category_key, name, sort_order, status, created_at, updated_at)
VALUES ($1, $2, $3, $4, $5, $6)
ON CONFLICT (category_key) DO UPDATE
SET name = EXCLUDED.name, sort_order = EXCLUDED.sort_order, status = EXCLUDED.status, updated_at = EXCLUDED.updated_at
RETURNING *
`,
[input.categoryKey, input.name, input.sortOrder, input.status, input.now, input.now],
)
return result.rows[0] || null
}
export async function countWorkCategoryUsages(categoryId: number | string): Promise<{
workOrderCount: number
productRuleCount: number
}> {
const [orders, rules] = await Promise.all([
query<{ total: number }>(
'SELECT COUNT(*)::int AS total FROM work_orders WHERE category_id = $1',
[Number(categoryId)],
),
query<{ total: number }>(
'SELECT COUNT(*)::int AS total FROM work_product_rules WHERE category_id = $1',
[Number(categoryId)],
),
])
return {
workOrderCount: Number(orders.rows[0]?.total || 0),
productRuleCount: Number(rules.rows[0]?.total || 0),
}
}
export async function deleteWorkCategory(
categoryId: number | string,
): Promise<{ deleted: boolean }> {
const result = await query<{ id: number }>(
'DELETE FROM work_categories WHERE id = $1 RETURNING id',
[Number(categoryId)],
)
return { deleted: Boolean(result.rows[0]) }
}
@@ -0,0 +1,82 @@
import type { PoolClient } from 'pg'
import { query } from '../../db/client.js'
export type WorkOrderEventRow = {
id: number
work_order_id: number
actor_type: string
actor_id: string
event_type: string
from_status: string
to_status: string
payload_json: string | Record<string, unknown>
created_at: string
}
export async function listWorkOrderEventsByOrderId(
workOrderId: number | string,
): Promise<WorkOrderEventRow[]> {
const result = await query<WorkOrderEventRow>(
`
SELECT *
FROM work_order_events
WHERE work_order_id = $1
ORDER BY id ASC
`,
[Number(workOrderId)],
)
return result.rows
}
export async function createWorkOrderEvent(input: {
workOrderId: number
actorType: string
actorId: string
eventType: string
fromStatus?: string
toStatus?: string
payloadJson?: string
now: string
}): Promise<void> {
await query(
`
INSERT INTO work_order_events (
work_order_id, actor_type, actor_id, event_type, from_status, to_status, payload_json, created_at
) VALUES ($1, $2, $3, $4, $5, $6, $7::jsonb, $8)
`,
[
input.workOrderId,
input.actorType,
input.actorId,
input.eventType,
input.fromStatus || '',
input.toStatus || '',
input.payloadJson || '{}',
input.now,
],
)
}
export async function createWorkOrderEventWithClient(
client: PoolClient,
input: Parameters<typeof createWorkOrderEvent>[0],
) {
await client.query(
`
INSERT INTO work_order_events (
work_order_id, actor_type, actor_id, event_type, from_status, to_status, payload_json, created_at
) VALUES ($1, $2, $3, $4, $5, $6, $7::jsonb, $8)
`,
[
input.workOrderId,
input.actorType,
input.actorId,
input.eventType,
input.fromStatus || '',
input.toStatus || '',
input.payloadJson || '{}',
input.now,
],
)
}
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,478 @@
import { query, withTransaction } from '../../db/client.js'
import {
ensureWorkerWalletWithClient,
getWorkerWalletWithClient,
toPositiveInteger,
WORK_ORDER_SHARE_SELECT,
} from './shared.js'
import { createWorkOrderEventWithClient } from './work-order-event-repo.js'
import {
countWorkerActiveOrdersWithClient,
getWorkOrderById,
isWorkOrderWithinHallCapacityWithClient,
} from './work-order-repo.js'
import type { WorkOrderRow, WorkOrderShareRow } from './types.js'
async function getWorkOrderShareById(shareId: number): Promise<WorkOrderShareRow | null> {
const result = await query<WorkOrderShareRow>(
`${WORK_ORDER_SHARE_SELECT} WHERE wos.id = $1 LIMIT 1`,
[shareId],
)
return result.rows[0] || null
}
export async function listWorkOrderShares(
workOrderId: number | string,
): Promise<WorkOrderShareRow[]> {
const result = await query<WorkOrderShareRow>(
`${WORK_ORDER_SHARE_SELECT} WHERE wos.work_order_id = $1 ORDER BY wos.id ASC`,
[Number(workOrderId)],
)
return result.rows
}
export async function listWorkOrderSharesByOrderIds(
workOrderIds: number[],
): Promise<WorkOrderShareRow[]> {
if (workOrderIds.length === 0) return []
const result = await query<WorkOrderShareRow>(
`${WORK_ORDER_SHARE_SELECT}
WHERE wos.work_order_id = ANY($1::bigint[])
ORDER BY wos.id ASC`,
[workOrderIds],
)
return result.rows
}
export async function listWorkerSharesByWorker(
workerId: number | string,
): Promise<WorkOrderShareRow[]> {
const result = await query<WorkOrderShareRow>(
`${WORK_ORDER_SHARE_SELECT}
WHERE wos.worker_id = $1 AND wos.status != 'cancelled'
ORDER BY wos.id DESC`,
[Number(workerId)],
)
return result.rows
}
export async function countWorkOrderPendingSharingSubmissions(
workOrderId: number | string,
): Promise<number> {
const result = await query<{ total: number }>(
`
SELECT COUNT(*)::int AS total
FROM work_order_shares
WHERE work_order_id = $1 AND status = 'joined'
`,
[Number(workOrderId)],
)
return Number(result.rows[0]?.total || 0)
}
export async function getWorkOrderShare(
workOrderId: number | string,
workerId: number | string,
): Promise<WorkOrderShareRow | null> {
const result = await query<WorkOrderShareRow>(
`${WORK_ORDER_SHARE_SELECT}
WHERE wos.work_order_id = $1 AND wos.worker_id = $2 AND wos.status != 'cancelled' LIMIT 1`,
[Number(workOrderId), Number(workerId)],
)
return result.rows[0] || null
}
/** 普通拼单使用指定份数,立即抢单则在事务中承接全部剩余份数。 */
export function resolveWorkOrderShareJoinQuantity(input: {
quantity: number
remaining: number
takeRemaining?: boolean
}): number {
return input.takeRemaining ? Math.max(0, input.remaining) : input.quantity
}
export async function joinWorkOrderShare(input: {
workOrderId: number
workerId: number
quantity: number
/** 立即抢单时在事务锁内承接全部剩余份数。 */
takeRemaining?: boolean
maxActiveOrders: number
hallCandidateLimit?: number
visibleAfterIso?: string
now: string
}): Promise<{
share: WorkOrderShareRow | null
failureReason:
| 'work_order_not_open'
| 'sharing_disabled'
| 'work_order_owner_conflict'
| 'sharing_quantity_full'
| 'worker_active_order_limit'
| 'worker_deposit_insufficient'
| null
}> {
return withTransaction(async (client) => {
await ensureWorkerWalletWithClient(client, input.workerId, input.now)
const orderResult = await client.query<WorkOrderRow>(
`
SELECT *
FROM work_orders
WHERE id = $1
FOR UPDATE
`,
[input.workOrderId],
)
const workOrder = orderResult.rows[0] || null
if (!workOrder || workOrder.status !== 'open') {
return { share: null, failureReason: 'work_order_not_open' }
}
if (
input.hallCandidateLimit &&
!(await isWorkOrderWithinHallCapacityWithClient(client, {
workOrderId: input.workOrderId,
hallCandidateLimit: input.hallCandidateLimit,
visibleAfterIso: input.visibleAfterIso,
}))
) {
return { share: null, failureReason: 'work_order_not_open' }
}
if (workOrder.sharing_enabled !== true) {
return { share: null, failureReason: 'sharing_disabled' }
}
if (Number(workOrder.assigned_worker_id || 0) === input.workerId) {
return { share: null, failureReason: 'work_order_owner_conflict' }
}
const totalQuantity = toPositiveInteger(workOrder.sharing_total_quantity, 1)
const joinedResult = await client.query<{ total: number }>(
`
SELECT COALESCE(SUM(quantity), 0)::int AS total
FROM work_order_shares
WHERE work_order_id = $1 AND status != 'cancelled'
`,
[input.workOrderId],
)
const joinedQuantity = Number(joinedResult.rows[0]?.total || 0)
const remaining = Math.max(0, totalQuantity - joinedQuantity)
const requestedQuantity = resolveWorkOrderShareJoinQuantity({
quantity: input.quantity,
remaining,
takeRemaining: input.takeRemaining === true,
})
if (requestedQuantity <= 0 || requestedQuantity > remaining) {
return { share: null, failureReason: 'sharing_quantity_full' }
}
const existingShareResult = await client.query<WorkOrderShareRow>(
`SELECT *
FROM work_order_shares
WHERE work_order_id = $1 AND worker_id = $2
FOR UPDATE`,
[input.workOrderId, input.workerId],
)
const existingShare = existingShareResult.rows[0] || null
if (existingShare && (!input.takeRemaining || existingShare.status !== 'joined')) {
return { share: null, failureReason: 'sharing_quantity_full' }
}
if (!existingShare) {
const activeOrderCount = await countWorkerActiveOrdersWithClient(client, input.workerId)
if (activeOrderCount >= input.maxActiveOrders) {
return { share: null, failureReason: 'worker_active_order_limit' }
}
}
const unitReward = toPositiveInteger(workOrder.sharing_unit_reward, 0)
const nextQuantity = Number(existingShare?.quantity || 0) + requestedQuantity
const shareReward = nextQuantity * unitReward
const nextShareDeposit = Math.round(
(Number(workOrder.required_deposit_amount || 0) * nextQuantity) / totalQuantity,
)
const depositToFreeze = Math.max(
0,
nextShareDeposit - Number(existingShare?.share_deposit || 0),
)
const wallet = await getWorkerWalletWithClient(client, input.workerId)
const available = Number(wallet?.available_amount || 0)
if (available < depositToFreeze) {
return { share: null, failureReason: 'worker_deposit_insufficient' }
}
const shareResult = existingShare
? await client.query<{ id: number }>(
`
UPDATE work_order_shares
SET
quantity = $1,
unit_reward = $2,
share_reward = $3,
share_deposit = $4,
updated_at = $5
WHERE id = $6
RETURNING id
`,
[nextQuantity, unitReward, shareReward, nextShareDeposit, input.now, existingShare.id],
)
: await client.query<{ id: number }>(
`
INSERT INTO work_order_shares (
work_order_id, worker_id, quantity, unit_reward, share_reward,
share_deposit, status, acceptance_json, created_at, updated_at
) VALUES ($1, $2, $3, $4, $5, $6, 'joined', '{}'::jsonb, $7, $7)
RETURNING id
`,
[
input.workOrderId,
input.workerId,
requestedQuantity,
unitReward,
shareReward,
nextShareDeposit,
input.now,
],
)
if (depositToFreeze > 0) {
const nextAvailable = available - depositToFreeze
const nextFrozen = Number(wallet?.frozen_deposit_amount || 0) + depositToFreeze
await client.query(
`
UPDATE worker_wallets
SET available_amount = $1, frozen_deposit_amount = $2, updated_at = $3
WHERE worker_id = $4
`,
[nextAvailable, nextFrozen, input.now, input.workerId],
)
await client.query(
`
INSERT INTO worker_wallet_ledgers (
worker_id, ledger_type, amount, balance_after, frozen_after,
related_work_order_id, note, payload_json, created_at
) VALUES ($1, 'deposit_freeze', $2, $3, $4, $5, $6, $7::jsonb, $8)
`,
[
input.workerId,
depositToFreeze,
nextAvailable,
nextFrozen,
input.workOrderId,
`拼单冻结押金 ${requestedQuantity}`,
JSON.stringify({ workOrderId: input.workOrderId, quantity: requestedQuantity }),
input.now,
],
)
}
await createWorkOrderEventWithClient(client, {
workOrderId: input.workOrderId,
actorType: 'worker',
actorId: String(input.workerId),
eventType: 'sharing_joined',
fromStatus: 'open',
toStatus: 'open',
payloadJson: JSON.stringify({
workerId: input.workerId,
quantity: requestedQuantity,
takeRemaining: input.takeRemaining === true,
}),
now: input.now,
})
const share = await client.query<WorkOrderShareRow>(
`${WORK_ORDER_SHARE_SELECT} WHERE wos.id = $1 LIMIT 1`,
[Number(shareResult.rows[0]?.id || 0)],
)
return { share: share.rows[0] || null, failureReason: null }
})
}
export async function submitWorkOrderShareAcceptance(input: {
workOrderId: number
workerId: number
acceptanceJson: string
now: string
}): Promise<{
share: WorkOrderShareRow | null
failureReason: 'share_not_found' | 'share_status_invalid' | null
}> {
return withTransaction(async (client) => {
const orderResult = await client.query<{
status: string
sharing_total_quantity: number
}>(
`SELECT status, sharing_total_quantity
FROM work_orders
WHERE id = $1
FOR UPDATE`,
[input.workOrderId],
)
const workOrder = orderResult.rows[0] || null
if (!workOrder || !['open', 'pending_acceptance'].includes(workOrder.status)) {
return { share: null, failureReason: 'share_not_found' }
}
const lockResult = await client.query<{ id: number }>(
`
SELECT id
FROM work_order_shares
WHERE work_order_id = $1 AND worker_id = $2
FOR UPDATE
`,
[input.workOrderId, input.workerId],
)
const shareId = lockResult.rows[0]?.id || 0
if (!shareId) {
return { share: null, failureReason: 'share_not_found' }
}
const shareResult = await client.query<WorkOrderShareRow>(
`${WORK_ORDER_SHARE_SELECT} WHERE wos.id = $1 LIMIT 1`,
[shareId],
)
const share = shareResult.rows[0] || null
if (!share) {
return { share: null, failureReason: 'share_not_found' }
}
if (share.status !== 'joined') {
return { share, failureReason: 'share_status_invalid' }
}
await client.query(
`
UPDATE work_order_shares
SET status = 'submitted', acceptance_json = $1::jsonb, draft_acceptance_json = '{}'::jsonb, submitted_at = $2, updated_at = $2
WHERE id = $3
`,
[input.acceptanceJson, input.now, Number(share.id)],
)
const progressResult = await client.query<{
joined_quantity: number
pending_submission_count: number
}>(
`SELECT
COALESCE(SUM(quantity), 0)::int AS joined_quantity,
COUNT(*) FILTER (WHERE status = 'joined')::int AS pending_submission_count
FROM work_order_shares
WHERE work_order_id = $1
AND status != 'cancelled'`,
[input.workOrderId],
)
const progress = progressResult.rows[0]
const allQuantityAllocated =
Number(progress?.joined_quantity || 0) >=
toPositiveInteger(workOrder.sharing_total_quantity, 1)
if (allQuantityAllocated && Number(progress?.pending_submission_count || 0) === 0) {
await client.query(
`UPDATE work_orders
SET
status = 'pending_acceptance',
submitted_at = COALESCE(submitted_at, $1),
updated_at = $1
WHERE id = $2
AND status = 'open'`,
[input.now, input.workOrderId],
)
}
await createWorkOrderEventWithClient(client, {
workOrderId: input.workOrderId,
actorType: 'worker',
actorId: String(input.workerId),
eventType: 'sharing_acceptance_submitted',
fromStatus: 'joined',
toStatus: 'submitted',
payloadJson: JSON.stringify({ workerId: input.workerId, quantity: share.quantity }),
now: input.now,
})
const updated = await client.query<WorkOrderShareRow>(
`${WORK_ORDER_SHARE_SELECT} WHERE wos.id = $1 LIMIT 1`,
[Number(share.id)],
)
return { share: updated.rows[0] || null, failureReason: null }
})
}
/** 保存整单验收草稿(暂存图片/说明),不改变工单状态 */
export async function updateWorkOrderDraftAcceptance(input: {
workOrderId: number
draftJson: string
now: string
}): Promise<WorkOrderRow | null> {
await query(
`
UPDATE work_orders
SET draft_acceptance_json = $1::jsonb, updated_at = $2
WHERE id = $3
`,
[input.draftJson, input.now, Number(input.workOrderId)],
)
return getWorkOrderById(input.workOrderId)
}
/** 更新整单已提交的验收资料(保留原 submitted_at,用于提交验收后补充/修改图片) */
export async function updateWorkOrderAcceptanceRecord(input: {
workOrderId: number
acceptanceJson: string
now: string
}): Promise<WorkOrderRow | null> {
await query(
`
UPDATE work_orders
SET acceptance_json = $1::jsonb, updated_at = $2
WHERE id = $3
`,
[input.acceptanceJson, input.now, Number(input.workOrderId)],
)
return getWorkOrderById(input.workOrderId)
}
/** 保存拼单份额的验收草稿(暂存图片/说明),不改变份额状态 */
export async function updateWorkOrderShareDraftAcceptance(input: {
workOrderId: number
workerId: number
draftJson: string
now: string
}): Promise<WorkOrderShareRow | null> {
const result = await query<WorkOrderShareRow>(
`${WORK_ORDER_SHARE_SELECT}
WHERE wos.work_order_id = $1 AND wos.worker_id = $2`,
[Number(input.workOrderId), Number(input.workerId)],
)
const share = result.rows[0] || null
if (!share) return null
await query(
`
UPDATE work_order_shares
SET draft_acceptance_json = $1::jsonb, updated_at = $2
WHERE id = $3
`,
[input.draftJson, input.now, Number(share.id)],
)
return getWorkOrderShareById(share.id)
}
/** 更新拼单份额已提交的验收资料(保留原 submitted_at,用于提交验收后补充/修改图片) */
export async function updateWorkOrderShareAcceptanceRecord(input: {
workOrderId: number
workerId: number
acceptanceJson: string
now: string
}): Promise<WorkOrderShareRow | null> {
const result = await query<WorkOrderShareRow>(
`${WORK_ORDER_SHARE_SELECT}
WHERE wos.work_order_id = $1 AND wos.worker_id = $2`,
[Number(input.workOrderId), Number(input.workerId)],
)
const share = result.rows[0] || null
if (!share) return null
await query(
`
UPDATE work_order_shares
SET acceptance_json = $1::jsonb, updated_at = $2
WHERE id = $3
`,
[input.acceptanceJson, input.now, Number(share.id)],
)
return getWorkOrderShareById(share.id)
}
@@ -0,0 +1,421 @@
import { query, withTransaction } from '../../db/client.js'
import { toPositiveInteger } from './shared.js'
import type {
ProductRuleListInput,
WorkOrderRow,
WorkOrderShareRow,
WorkProductRuleRow,
} from './types.js'
import { WORK_ORDER_SELECT } from './work-order-repo.js'
const WORK_PRODUCT_RULE_SELECT = `
SELECT
wpr.*,
wc.name AS category_name
FROM work_product_rules wpr
LEFT JOIN work_categories wc ON wc.id = wpr.category_id
`
function buildWorkProductRuleWhere({
enabled = null,
keyword = '',
categoryId = 0,
}: ProductRuleListInput) {
const filters: string[] = []
const params: unknown[] = []
if (enabled !== null && enabled !== undefined) {
params.push(Boolean(enabled))
filters.push(`wpr.enabled = $${params.length}`)
}
if (keyword) {
params.push(`%${keyword}%`)
filters.push(`(
wpr.rule_key ILIKE $${params.length}
OR wpr.sku_code ILIKE $${params.length}
OR wpr.product_name ILIKE $${params.length}
)`)
}
if (categoryId) {
params.push(categoryId)
filters.push(`wpr.category_id = $${params.length}`)
}
return {
whereClause: filters.length > 0 ? `WHERE ${filters.join(' AND ')}` : '',
params,
}
}
export async function listWorkProductRules({
enabled = null,
keyword = '',
categoryId = 0,
}: ProductRuleListInput = {}): Promise<WorkProductRuleRow[]> {
const { whereClause, params } = buildWorkProductRuleWhere({ enabled, keyword, categoryId })
const result = await query<WorkProductRuleRow>(
`${WORK_PRODUCT_RULE_SELECT}
${whereClause}
ORDER BY wpr.sort_order ASC, wpr.id DESC`,
params,
)
return result.rows
}
export async function listWorkProductRulesPage({
page = 1,
pageSize = 20,
enabled = null,
keyword = '',
categoryId = 0,
}: ProductRuleListInput): Promise<{ items: WorkProductRuleRow[]; total: number }> {
const { whereClause, params } = buildWorkProductRuleWhere({ enabled, keyword, categoryId })
const totalResult = await query<{ total: number }>(
`
SELECT COUNT(*)::int AS total
FROM work_product_rules wpr
${whereClause}
`,
params,
)
const offset = (page - 1) * pageSize
const listParams = [...params, pageSize, offset]
const result = await query<WorkProductRuleRow>(
`${WORK_PRODUCT_RULE_SELECT}
${whereClause}
ORDER BY wpr.sort_order ASC, wpr.id DESC
LIMIT $${listParams.length - 1} OFFSET $${listParams.length}`,
listParams,
)
return { items: result.rows, total: Number(totalResult.rows[0]?.total || 0) }
}
export async function countWorkProductRulesByCategory(
input: {
enabled?: boolean | null
keyword?: string
} = {},
): Promise<Array<{ categoryId: number | null; total: number }>> {
const { whereClause, params } = buildWorkProductRuleWhere(input)
const result = await query<{ category_id: number | null; total: number }>(
`
SELECT wpr.category_id, COUNT(*)::int AS total
FROM work_product_rules wpr
${whereClause}
GROUP BY wpr.category_id
`,
params,
)
return result.rows.map((row) => ({
categoryId: row.category_id ? Number(row.category_id) : null,
total: Number(row.total || 0),
}))
}
export async function getWorkProductRuleByKey(ruleKey: string): Promise<WorkProductRuleRow | null> {
const result = await query<WorkProductRuleRow>(
`${WORK_PRODUCT_RULE_SELECT} WHERE wpr.rule_key = $1 LIMIT 1`,
[String(ruleKey || '').trim()],
)
return result.rows[0] || null
}
export async function getWorkProductRuleById(
ruleId: number | string,
): Promise<WorkProductRuleRow | null> {
const result = await query<WorkProductRuleRow>(
`${WORK_PRODUCT_RULE_SELECT} WHERE wpr.id = $1 LIMIT 1`,
[Number(ruleId)],
)
return result.rows[0] || null
}
/** 查询仍在接单大厅、且明确来源于指定模板的工单。 */
export async function listOpenWorkOrdersByProductRuleId(
ruleId: number | string,
): Promise<WorkOrderRow[]> {
const result = await query<WorkOrderRow>(
`${WORK_ORDER_SELECT}
WHERE wo.product_rule_id = $1
AND wo.status = 'open'
ORDER BY wo.id ASC`,
[Number(ruleId)],
)
return result.rows
}
export type SyncOpenWorkOrderPriceResult = {
updated: boolean
reason: 'updated' | 'price_unchanged' | 'not_open' | 'assigned' | 'sharing_filled'
workOrderId: number
sharingPartial: boolean
workerIds: number[]
}
/** 拼单改价时保留已加入份额的报酬,只按新价计算剩余份额。 */
export function resolveOpenSharingSyncReward(input: {
totalQuantity: number
joinedQuantity: number
lockedShareReward: number
nextUnitReward: number
}) {
const totalQuantity = toPositiveInteger(input.totalQuantity, 1)
const joinedQuantity = Math.min(totalQuantity, Math.max(0, Number(input.joinedQuantity || 0)))
return (
Math.max(0, Number(input.lockedShareReward || 0)) +
(totalQuantity - joinedQuantity) * toPositiveInteger(input.nextUnitReward, 0)
)
}
/**
* 同步大厅工单的模板价格。
* 已加入的拼单份额报酬是合同快照,只更新尚未加入的剩余份额价格。
*/
export async function syncOpenWorkOrderPriceFromRule(input: {
workOrderId: number
productRuleId: number
rewardAmount: number
sharingUnitReward?: number
actorName: string
now: string
}): Promise<SyncOpenWorkOrderPriceResult> {
return withTransaction(async (client) => {
const orderResult = await client.query<WorkOrderRow>(
`SELECT *
FROM work_orders
WHERE id = $1 AND product_rule_id = $2
FOR UPDATE`,
[input.workOrderId, input.productRuleId],
)
const workOrder = orderResult.rows[0] || null
if (!workOrder || workOrder.status !== 'open') {
return {
updated: false,
reason: 'not_open',
workOrderId: input.workOrderId,
sharingPartial: false,
workerIds: [],
}
}
if (workOrder.assigned_worker_id) {
return {
updated: false,
reason: 'assigned',
workOrderId: input.workOrderId,
sharingPartial: false,
workerIds: [],
}
}
const sharesResult = await client.query<WorkOrderShareRow>(
`SELECT *
FROM work_order_shares
WHERE work_order_id = $1
FOR UPDATE`,
[workOrder.id],
)
const activeShares = sharesResult.rows.filter((share) => share.status !== 'cancelled')
const joinedQuantity = activeShares.reduce(
(sum, share) => sum + Math.max(0, Number(share.quantity || 0)),
0,
)
const workerIds = [
...new Set(
activeShares.map((share) => Number(share.worker_id)).filter((workerId) => workerId > 0),
),
]
const sharingPartial = workOrder.sharing_enabled === true && joinedQuantity > 0
let nextRewardAmount = toPositiveInteger(input.rewardAmount, 0)
let nextSharingUnitReward = Number(workOrder.sharing_unit_reward || 0)
if (workOrder.sharing_enabled === true) {
const totalQuantity = toPositiveInteger(workOrder.sharing_total_quantity, 1)
if (joinedQuantity >= totalQuantity) {
return {
updated: false,
reason: 'sharing_filled',
workOrderId: Number(workOrder.id),
sharingPartial,
workerIds,
}
}
nextSharingUnitReward = toPositiveInteger(input.sharingUnitReward, 0)
const lockedShareReward = activeShares.reduce(
(sum, share) => sum + Math.max(0, Number(share.share_reward || 0)),
0,
)
nextRewardAmount = resolveOpenSharingSyncReward({
totalQuantity,
joinedQuantity,
lockedShareReward,
nextUnitReward: nextSharingUnitReward,
})
}
const unchanged =
Number(workOrder.reward_amount || 0) === nextRewardAmount &&
(!workOrder.sharing_enabled ||
Number(workOrder.sharing_unit_reward || 0) === nextSharingUnitReward)
if (unchanged) {
return {
updated: false,
reason: 'price_unchanged',
workOrderId: Number(workOrder.id),
sharingPartial,
workerIds,
}
}
await client.query(
`UPDATE work_orders
SET reward_amount = $1,
sharing_unit_reward = $2,
updated_at = $3
WHERE id = $4`,
[nextRewardAmount, nextSharingUnitReward, input.now, workOrder.id],
)
await client.query(
`INSERT INTO work_order_events (
work_order_id, actor_type, actor_id, event_type, from_status, to_status, payload_json, created_at
) VALUES ($1, 'admin', $2, 'price_synced_from_rule', $3, $3, $4::jsonb, $5)`,
[
workOrder.id,
input.actorName,
workOrder.status,
JSON.stringify({
productRuleId: input.productRuleId,
previousRewardAmount: Number(workOrder.reward_amount || 0),
rewardAmount: nextRewardAmount,
previousSharingUnitReward: Number(workOrder.sharing_unit_reward || 0),
sharingUnitReward: nextSharingUnitReward,
joinedQuantity,
lockedShareReward: activeShares.reduce(
(sum, share) => sum + Math.max(0, Number(share.share_reward || 0)),
0,
),
}),
input.now,
],
)
return {
updated: true,
reason: 'updated',
workOrderId: Number(workOrder.id),
sharingPartial,
workerIds,
}
})
}
export async function deleteWorkProductRule(
ruleId: number | string,
): Promise<{ deleted: boolean }> {
const result = await query<{ id: number }>(
'DELETE FROM work_product_rules WHERE id = $1 RETURNING id',
[Number(ruleId)],
)
return { deleted: Boolean(result.rows[0]) }
}
export async function upsertWorkProductRule(input: {
ruleKey: string
provider: string
platform: string
shopId: string
skuCode: string
productName: string
matchType: string
categoryId: number | null
enabled: boolean
autoCreate: boolean
rewardAmount: number
unitPriceFen?: number
requiredDepositAmount: number
depositThresholdAmount: number
sharingEnabled?: boolean
sharingTotalQuantity?: number
sharingUnitReward?: number
sharingAutoFromOrder?: boolean
timeoutMinutes?: number
timeoutPolicy?: string
requirementJson: string
matchJson?: string
sortOrder: number
now: string
}): Promise<WorkProductRuleRow | null> {
const result = await query<WorkProductRuleRow>(
`
INSERT INTO work_product_rules (
rule_key, provider, platform, shop_id, sku_code, product_name,
match_type, category_id, enabled, auto_create, reward_amount,
unit_price_fen,
required_deposit_amount, deposit_threshold_amount,
sharing_enabled, sharing_total_quantity, sharing_unit_reward, sharing_auto_from_order,
timeout_minutes, timeout_policy,
requirement_json, match_json,
sort_order, created_at, updated_at
) VALUES (
$1, $2, $3, $4, $5, $6,
$7, $8, $9, $10, $11,
$12,
$13, $14,
$15, $16, $17, $18,
$19, $20,
$21::jsonb, $22::jsonb,
$23, $24, $25
)
ON CONFLICT (rule_key) DO UPDATE
SET
provider = EXCLUDED.provider,
platform = EXCLUDED.platform,
shop_id = EXCLUDED.shop_id,
sku_code = EXCLUDED.sku_code,
product_name = EXCLUDED.product_name,
match_type = EXCLUDED.match_type,
category_id = EXCLUDED.category_id,
enabled = EXCLUDED.enabled,
auto_create = EXCLUDED.auto_create,
reward_amount = EXCLUDED.reward_amount,
unit_price_fen = EXCLUDED.unit_price_fen,
required_deposit_amount = EXCLUDED.required_deposit_amount,
deposit_threshold_amount = EXCLUDED.deposit_threshold_amount,
sharing_enabled = EXCLUDED.sharing_enabled,
sharing_total_quantity = EXCLUDED.sharing_total_quantity,
sharing_unit_reward = EXCLUDED.sharing_unit_reward,
sharing_auto_from_order = EXCLUDED.sharing_auto_from_order,
timeout_minutes = EXCLUDED.timeout_minutes,
timeout_policy = EXCLUDED.timeout_policy,
requirement_json = EXCLUDED.requirement_json,
match_json = EXCLUDED.match_json,
sort_order = EXCLUDED.sort_order,
updated_at = EXCLUDED.updated_at
RETURNING *
`,
[
input.ruleKey,
input.provider,
input.platform,
input.shopId,
input.skuCode,
input.productName,
input.matchType,
input.categoryId,
input.enabled,
input.autoCreate,
input.rewardAmount,
Math.max(0, toPositiveInteger(input.unitPriceFen, 0)),
input.requiredDepositAmount,
input.depositThresholdAmount,
input.sharingEnabled === true,
toPositiveInteger(input.sharingTotalQuantity, 1),
toPositiveInteger(input.sharingUnitReward, 0),
input.sharingAutoFromOrder === true,
toPositiveInteger(input.timeoutMinutes, 0),
String(input.timeoutPolicy || 'reopen').trim(),
input.requirementJson,
input.matchJson || '{}',
input.sortOrder,
input.now,
input.now,
],
)
return getWorkProductRuleByKey(result.rows[0]?.rule_key || input.ruleKey)
}
@@ -0,0 +1,609 @@
import type { PoolClient } from 'pg'
import { query, withTransaction } from '../../db/client.js'
import {
ensureWorkerWalletWithClient,
getWorkerWalletWithClient,
resolveWorkOrderShareCancellationStatus,
} from './shared.js'
import { createWorkOrderEventWithClient } from './work-order-event-repo.js'
import {
getOutstandingDepositAmountWithClient,
getWorkOrderByIdWithClient,
} from './work-order-repo.js'
import type { WorkerCancelRequestRow, WorkOrderRow, WorkOrderShareRow } from './types.js'
const WORKER_CANCEL_REQUEST_SELECT = `
SELECT
cr.*,
COALESCE(NULLIF(BTRIM(wo.platform_order_id), ''), o.platform_order_id, ft.platform_order_id, '') AS platform_order_id,
COALESCE(
NULLIF(BTRIM(o.shop_name), ''),
NULLIF(BTRIM(o.shop_id), ''),
NULLIF(BTRIM(ft.shop_name), ''),
NULLIF(BTRIM(ft.shop_id), ''),
''
) AS shop_name,
wo.product_name,
wu.username AS worker_username,
wu.display_name AS worker_display_name,
wos.quantity AS share_quantity
FROM worker_order_cancel_requests cr
INNER JOIN work_orders wo ON wo.id = cr.work_order_id
LEFT JOIN orders o ON o.id = wo.order_id
LEFT JOIN fulfillment_tasks ft ON ft.id = wo.task_id
INNER JOIN worker_users wu ON wu.id = cr.worker_id
LEFT JOIN work_order_shares wos ON wos.id = cr.work_order_share_id
`
async function getWorkerCancelRequestByIdWithClient(
client: PoolClient,
requestId: number | string,
): Promise<WorkerCancelRequestRow | null> {
const result = await client.query<WorkerCancelRequestRow>(
`${WORKER_CANCEL_REQUEST_SELECT} WHERE cr.id = $1 LIMIT 1`,
[Number(requestId)],
)
return result.rows[0] || null
}
async function countWorkerAcceptedOrdersWithClient(
client: PoolClient,
workerId: number,
): Promise<number> {
const result = await client.query<{ total: number }>(
`SELECT COUNT(DISTINCT order_key)::int AS total
FROM (
SELECT COALESCE(
'order:' || NULLIF(wo.order_id::text, ''),
'platform_order:' || NULLIF(BTRIM(wo.platform_order_id), ''),
'work_order:' || wo.id::text
) AS order_key
FROM work_orders wo
WHERE wo.assigned_worker_id = $1 AND wo.status = 'accepted'
UNION
SELECT COALESCE(
'order:' || NULLIF(wo.order_id::text, ''),
'platform_order:' || NULLIF(BTRIM(wo.platform_order_id), ''),
'work_order:' || wo.id::text
) AS order_key
FROM work_order_shares wos
INNER JOIN work_orders wo ON wo.id = wos.work_order_id
WHERE wos.worker_id = $1 AND wos.status = 'accepted'
) accepted_orders`,
[workerId],
)
return Number(result.rows[0]?.total || 0)
}
async function applyWorkerCancelRequestPenaltyWithClient(
client: PoolClient,
input: { workerId: number; requestId: number; now: string },
): Promise<{ levelBeforeId: number | null; levelChanged: boolean; accountDisabled: boolean }> {
const workerResult = await client.query<{
level_id: number | null
level_key: string
level_name: string
sort_order: number
level_penalty_accepted_count: number
}>(
`SELECT wu.level_id, wl.level_key, wl.name AS level_name, wl.sort_order,
wu.level_penalty_accepted_count
FROM worker_users wu
LEFT JOIN worker_levels wl ON wl.id = wu.level_id
WHERE wu.id = $1
FOR UPDATE OF wu`,
[input.workerId],
)
const worker = workerResult.rows[0]
const levelBeforeId = worker?.level_id ? Number(worker.level_id) : null
if (!worker || !levelBeforeId) {
return { levelBeforeId, levelChanged: false, accountDisabled: false }
}
const approvedResult = await client.query<{ total: number }>(
`SELECT COUNT(*)::int AS total
FROM worker_order_cancel_requests
WHERE worker_id = $1 AND status = 'approved'`,
[input.workerId],
)
const approvedTotal = Number(approvedResult.rows[0]?.total || 0) + 1
const isVip1 = String(worker.level_key || '').toLowerCase() === 'vip1'
if (isVip1) {
const vip1ApprovedResult = await client.query<{ total: number }>(
`SELECT COUNT(*)::int AS total
FROM worker_order_cancel_requests
WHERE worker_id = $1 AND status = 'approved' AND level_before_id = $2`,
[input.workerId, levelBeforeId],
)
const vip1ApprovedTotal = Number(vip1ApprovedResult.rows[0]?.total || 0) + 1
if (vip1ApprovedTotal >= 5) {
await client.query(
`UPDATE worker_users
SET status = 'disabled', session_version = session_version + 1, updated_at = $1
WHERE id = $2`,
[input.now, input.workerId],
)
await client.query(
`UPDATE worker_sessions
SET status = 'revoked', revoked_at = $1, updated_at = $1
WHERE worker_id = $2 AND status = 'active'`,
[input.now, input.workerId],
)
await client.query(
`INSERT INTO worker_level_logs (
worker_id, from_level_id, to_level_id, reason, payload_json, created_at
) VALUES ($1, $2, $2, 'cancel_request_ban', $3::jsonb, $4)`,
[
input.workerId,
levelBeforeId,
JSON.stringify({ requestId: input.requestId, vip1ApprovedTotal, approvedTotal }),
input.now,
],
)
return { levelBeforeId, levelChanged: false, accountDisabled: true }
}
return { levelBeforeId, levelChanged: false, accountDisabled: false }
}
if (approvedTotal % 10 !== 0) {
return { levelBeforeId, levelChanged: false, accountDisabled: false }
}
const lowerLevelResult = await client.query<{
id: number
level_key: string
threshold: number
}>(
`SELECT id, level_key, COALESCE((permission_json->>'upgradeThreshold')::int, 0) AS threshold
FROM worker_levels
WHERE status = 'active' AND sort_order < $1
ORDER BY sort_order DESC, id DESC
LIMIT 1`,
[Number(worker.sort_order || 0)],
)
const lowerLevel = lowerLevelResult.rows[0]
if (!lowerLevel) {
return { levelBeforeId, levelChanged: false, accountDisabled: false }
}
const actualAcceptedCount = await countWorkerAcceptedOrdersWithClient(client, input.workerId)
const nextPenaltyAcceptedCount = Math.max(
Number(worker.level_penalty_accepted_count || 0),
actualAcceptedCount - Math.max(0, Number(lowerLevel.threshold || 0)),
)
await client.query(
`UPDATE worker_users
SET level_id = $1, level_penalty_accepted_count = $2, updated_at = $3
WHERE id = $4`,
[lowerLevel.id, nextPenaltyAcceptedCount, input.now, input.workerId],
)
await client.query(
`INSERT INTO worker_level_logs (
worker_id, from_level_id, to_level_id, reason, payload_json, created_at
) VALUES ($1, $2, $3, 'cancel_request_penalty', $4::jsonb, $5)`,
[
input.workerId,
levelBeforeId,
lowerLevel.id,
JSON.stringify({
requestId: input.requestId,
approvedTotal,
actualAcceptedCount,
penaltyAcceptedCount: nextPenaltyAcceptedCount,
}),
input.now,
],
)
return { levelBeforeId, levelChanged: true, accountDisabled: false }
}
/** 打手提交撤单申请;待审核期间订单与押金均保持原状。 */
export async function createWorkerCancelRequest(input: {
workOrderId: number
workerId: number
reason: string
now: string
}): Promise<{
request: WorkerCancelRequestRow | null
dailyUsed: number
failureReason:
| 'work_order_not_cancellable'
| 'work_order_owner_required'
| 'share_not_cancellable'
| 'request_pending'
| 'daily_limit_reached'
| null
}> {
return withTransaction(async (client) => {
const orderResult = await client.query<WorkOrderRow>(
'SELECT * FROM work_orders WHERE id = $1 FOR UPDATE',
[input.workOrderId],
)
const workOrder = orderResult.rows[0] || null
if (!workOrder) {
return { request: null, dailyUsed: 0, failureReason: 'work_order_not_cancellable' }
}
let shareId: number | null = null
if (workOrder.sharing_enabled === true) {
const shareResult = await client.query<WorkOrderShareRow>(
`SELECT * FROM work_order_shares
WHERE work_order_id = $1 AND worker_id = $2
FOR UPDATE`,
[input.workOrderId, input.workerId],
)
const share = shareResult.rows[0] || null
if (!share || !resolveWorkOrderShareCancellationStatus(workOrder.status, share.status)) {
return { request: null, dailyUsed: 0, failureReason: 'share_not_cancellable' }
}
shareId = Number(share.id)
} else if (
workOrder.status !== 'in_progress' ||
Number(workOrder.assigned_worker_id || 0) !== input.workerId
) {
return {
request: null,
dailyUsed: 0,
failureReason: workOrder.assigned_worker_id
? 'work_order_not_cancellable'
: 'work_order_owner_required',
}
}
const dailyResult = await client.query<{ total: number }>(
`SELECT COUNT(*)::int AS total
FROM worker_order_cancel_requests
WHERE worker_id = $1
AND status IN ('pending', 'approved')
AND created_at >= date_trunc('day', $2::timestamptz AT TIME ZONE 'Asia/Shanghai') AT TIME ZONE 'Asia/Shanghai'`,
[input.workerId, input.now],
)
const dailyUsed = Number(dailyResult.rows[0]?.total || 0)
if (dailyUsed >= 5) {
return { request: null, dailyUsed, failureReason: 'daily_limit_reached' }
}
const pendingResult = await client.query<{ id: number }>(
`SELECT id FROM worker_order_cancel_requests
WHERE work_order_id = $1
AND worker_id = $2
AND COALESCE(work_order_share_id, 0) = $3
AND status = 'pending'
LIMIT 1`,
[input.workOrderId, input.workerId, shareId || 0],
)
if (pendingResult.rows[0]) {
return { request: null, dailyUsed, failureReason: 'request_pending' }
}
const created = await client.query<{ id: number }>(
`INSERT INTO worker_order_cancel_requests (
work_order_id, work_order_share_id, worker_id, reason, status, created_at, updated_at
) VALUES ($1, $2, $3, $4, 'pending', $5, $5)
RETURNING id`,
[input.workOrderId, shareId, input.workerId, input.reason, input.now],
)
const requestId = Number(created.rows[0]?.id || 0)
await createWorkOrderEventWithClient(client, {
workOrderId: input.workOrderId,
actorType: 'worker',
actorId: String(input.workerId),
eventType: 'cancel_request_submitted',
fromStatus: workOrder.status,
toStatus: workOrder.status,
payloadJson: JSON.stringify({ requestId, shareId, reason: input.reason }),
now: input.now,
})
return {
request: await getWorkerCancelRequestByIdWithClient(client, requestId),
dailyUsed: dailyUsed + 1,
failureReason: null,
}
})
}
export async function listWorkerCancelRequests(
workerId: number | string,
): Promise<WorkerCancelRequestRow[]> {
const result = await query<WorkerCancelRequestRow>(
`${WORKER_CANCEL_REQUEST_SELECT}
WHERE cr.worker_id = $1
ORDER BY cr.id DESC
LIMIT 100`,
[Number(workerId)],
)
return result.rows
}
export async function listAdminWorkerCancelRequests(
status = 'pending',
): Promise<WorkerCancelRequestRow[]> {
const result = await query<WorkerCancelRequestRow>(
`${WORKER_CANCEL_REQUEST_SELECT}
WHERE ($1 = '' OR cr.status = $1)
ORDER BY cr.id DESC
LIMIT 200`,
[status],
)
return result.rows
}
export async function getWorkerCancelRequestRisk(workerId: number | string, now: string) {
const result = await query<{ daily_used: number; approved_total: number }>(
`SELECT
COUNT(*) FILTER (
WHERE status IN ('pending', 'approved')
AND created_at >= date_trunc('day', $2::timestamptz AT TIME ZONE 'Asia/Shanghai') AT TIME ZONE 'Asia/Shanghai'
)::int AS daily_used,
COUNT(*) FILTER (WHERE status = 'approved')::int AS approved_total
FROM worker_order_cancel_requests
WHERE worker_id = $1`,
[Number(workerId), now],
)
const row = result.rows[0]
const approvedTotal = Number(row?.approved_total || 0)
const nextPenaltyAt = Math.ceil((approvedTotal + 1) / 10) * 10
return {
dailyUsed: Number(row?.daily_used || 0),
dailyLimit: 5,
approvedTotal,
nextPenaltyAt,
remainingToPenalty: Math.max(0, nextPenaltyAt - approvedTotal),
}
}
/** 客服审核通过时在同一事务内撤单、退押金并执行等级处罚。 */
export async function reviewWorkerCancelRequest(input: {
requestId: number
approved: boolean
reviewNote: string
actorName: string
now: string
}): Promise<{
request: WorkerCancelRequestRow | null
order: WorkOrderRow | null
releasedDepositAmount: number
levelChanged: boolean
accountDisabled: boolean
failureReason: 'request_not_found' | 'request_not_pending' | 'target_state_changed' | null
}> {
return withTransaction(async (client) => {
const requestResult = await client.query<WorkerCancelRequestRow>(
'SELECT * FROM worker_order_cancel_requests WHERE id = $1 FOR UPDATE',
[input.requestId],
)
const request = requestResult.rows[0] || null
if (!request) {
return {
request: null,
order: null,
releasedDepositAmount: 0,
levelChanged: false,
accountDisabled: false,
failureReason: 'request_not_found',
}
}
if (request.status !== 'pending') {
return {
request: await getWorkerCancelRequestByIdWithClient(client, request.id),
order: null,
releasedDepositAmount: 0,
levelChanged: false,
accountDisabled: false,
failureReason: 'request_not_pending',
}
}
const workOrderResult = await client.query<WorkOrderRow>(
'SELECT * FROM work_orders WHERE id = $1 FOR UPDATE',
[request.work_order_id],
)
const workOrder = workOrderResult.rows[0] || null
if (!workOrder) {
return {
request: null,
order: null,
releasedDepositAmount: 0,
levelChanged: false,
accountDisabled: false,
failureReason: 'target_state_changed',
}
}
if (!input.approved) {
await client.query(
`UPDATE worker_order_cancel_requests
SET status = 'rejected', review_note = $1, reviewed_by = $2, reviewed_at = $3, updated_at = $3
WHERE id = $4`,
[input.reviewNote, input.actorName, input.now, request.id],
)
await createWorkOrderEventWithClient(client, {
workOrderId: Number(workOrder.id),
actorType: 'admin',
actorId: input.actorName,
eventType: 'cancel_request_rejected',
fromStatus: workOrder.status,
toStatus: workOrder.status,
payloadJson: JSON.stringify({ requestId: request.id, note: input.reviewNote }),
now: input.now,
})
return {
request: await getWorkerCancelRequestByIdWithClient(client, request.id),
order: await getWorkOrderByIdWithClient(client, workOrder.id),
releasedDepositAmount: 0,
levelChanged: false,
accountDisabled: false,
failureReason: null,
}
}
let releasedDepositAmount = 0
if (request.work_order_share_id) {
const shareResult = await client.query<WorkOrderShareRow>(
`SELECT * FROM work_order_shares WHERE id = $1 AND work_order_id = $2 FOR UPDATE`,
[request.work_order_share_id, request.work_order_id],
)
const share = shareResult.rows[0] || null
if (
!share ||
Number(share.worker_id) !== Number(request.worker_id) ||
!resolveWorkOrderShareCancellationStatus(workOrder.status, share.status)
) {
return {
request: null,
order: null,
releasedDepositAmount: 0,
levelChanged: false,
accountDisabled: false,
failureReason: 'target_state_changed',
}
}
await ensureWorkerWalletWithClient(client, Number(request.worker_id), input.now)
const wallet = await getWorkerWalletWithClient(client, Number(request.worker_id))
releasedDepositAmount = Math.min(
Number(share.share_deposit || 0),
Number(wallet?.frozen_deposit_amount || 0),
)
if (releasedDepositAmount > 0) {
const nextAvailable = Number(wallet?.available_amount || 0) + releasedDepositAmount
const nextFrozen = Math.max(
0,
Number(wallet?.frozen_deposit_amount || 0) - releasedDepositAmount,
)
await client.query(
`UPDATE worker_wallets
SET available_amount = $1, frozen_deposit_amount = $2, updated_at = $3
WHERE worker_id = $4`,
[nextAvailable, nextFrozen, input.now, request.worker_id],
)
await client.query(
`INSERT INTO worker_wallet_ledgers (
worker_id, ledger_type, amount, balance_after, frozen_after,
related_work_order_id, note, payload_json, created_at
) VALUES ($1, 'deposit_release', $2, $3, $4, $5, '撤单申请审核通过退还拼单押金', $6::jsonb, $7)`,
[
request.worker_id,
releasedDepositAmount,
nextAvailable,
nextFrozen,
request.work_order_id,
JSON.stringify({ requestId: request.id, shareId: share.id, reason: request.reason }),
input.now,
],
)
}
await client.query(
`UPDATE work_order_shares SET status = 'cancelled', updated_at = $1 WHERE id = $2`,
[input.now, share.id],
)
await client.query(
`UPDATE work_orders
SET status = 'open', submitted_at = NULL, published_at = COALESCE(published_at, $1),
hall_queued_at = COALESCE(hall_queued_at, $1), updated_at = $1
WHERE id = $2`,
[input.now, workOrder.id],
)
} else {
if (
workOrder.status !== 'in_progress' ||
Number(workOrder.assigned_worker_id || 0) !== Number(request.worker_id)
) {
return {
request: null,
order: null,
releasedDepositAmount: 0,
levelChanged: false,
accountDisabled: false,
failureReason: 'target_state_changed',
}
}
await ensureWorkerWalletWithClient(client, Number(request.worker_id), input.now)
const wallet = await getWorkerWalletWithClient(client, Number(request.worker_id))
releasedDepositAmount = Math.min(
await getOutstandingDepositAmountWithClient(
client,
Number(request.worker_id),
workOrder.id,
),
Number(wallet?.frozen_deposit_amount || 0),
)
await client.query(
`UPDATE work_orders
SET status = 'open', assigned_worker_id = NULL, assigned_at = NULL, deadline_at = NULL,
published_at = $1, hall_queued_at = $1, updated_at = $1
WHERE id = $2`,
[input.now, workOrder.id],
)
if (releasedDepositAmount > 0) {
const nextAvailable = Number(wallet?.available_amount || 0) + releasedDepositAmount
const nextFrozen = Math.max(
0,
Number(wallet?.frozen_deposit_amount || 0) - releasedDepositAmount,
)
await client.query(
`UPDATE worker_wallets
SET available_amount = $1, frozen_deposit_amount = $2, updated_at = $3
WHERE worker_id = $4`,
[nextAvailable, nextFrozen, input.now, request.worker_id],
)
await client.query(
`INSERT INTO worker_wallet_ledgers (
worker_id, ledger_type, amount, balance_after, frozen_after,
related_work_order_id, note, payload_json, created_at
) VALUES ($1, 'deposit_release', $2, $3, $4, $5, '撤单申请审核通过退还押金', $6::jsonb, $7)`,
[
request.worker_id,
releasedDepositAmount,
nextAvailable,
nextFrozen,
request.work_order_id,
JSON.stringify({ requestId: request.id, reason: request.reason }),
input.now,
],
)
}
}
const penalty = await applyWorkerCancelRequestPenaltyWithClient(client, {
workerId: Number(request.worker_id),
requestId: Number(request.id),
now: input.now,
})
await client.query(
`UPDATE worker_order_cancel_requests
SET status = 'approved', review_note = $1, reviewed_by = $2, reviewed_at = $3,
level_before_id = $4, updated_at = $3
WHERE id = $5`,
[input.reviewNote, input.actorName, input.now, penalty.levelBeforeId, request.id],
)
const updatedOrder = await getWorkOrderByIdWithClient(client, workOrder.id)
await createWorkOrderEventWithClient(client, {
workOrderId: Number(workOrder.id),
actorType: 'admin',
actorId: input.actorName,
eventType: 'cancel_request_approved',
fromStatus: workOrder.status,
toStatus: updatedOrder?.status || workOrder.status,
payloadJson: JSON.stringify({
requestId: request.id,
shareId: request.work_order_share_id,
releasedDepositAmount,
levelChanged: penalty.levelChanged,
accountDisabled: penalty.accountDisabled,
note: input.reviewNote,
}),
now: input.now,
})
return {
request: await getWorkerCancelRequestByIdWithClient(client, request.id),
order: updatedOrder,
releasedDepositAmount,
levelChanged: penalty.levelChanged,
accountDisabled: penalty.accountDisabled,
failureReason: null,
}
})
}
@@ -0,0 +1,175 @@
import type { PoolClient } from 'pg'
import { query, withTransaction } from '../../db/client.js'
import { createWorkOrderEventWithClient } from './work-order-event-repo.js'
import { getWorkOrderByIdWithClient } from './work-order-repo.js'
import type { WorkerOrderFeedbackRow, WorkOrderRow, WorkOrderShareRow } from './types.js'
const WORKER_ORDER_FEEDBACK_SELECT = `
SELECT
wf.*,
COALESCE(NULLIF(BTRIM(wo.platform_order_id), ''), wo.work_order_no) AS work_order_no,
wo.product_name,
COALESCE(
NULLIF(BTRIM(o.shop_name), ''),
NULLIF(BTRIM(o.shop_id), ''),
NULLIF(BTRIM(ft.shop_name), ''),
NULLIF(BTRIM(ft.shop_id), ''),
''
) AS shop_name,
wu.username AS worker_username,
wu.display_name AS worker_display_name,
wos.quantity AS share_quantity
FROM worker_order_feedbacks wf
INNER JOIN work_orders wo ON wo.id = wf.work_order_id
LEFT JOIN orders o ON o.id = wo.order_id
LEFT JOIN fulfillment_tasks ft ON ft.id = wo.task_id
INNER JOIN worker_users wu ON wu.id = wf.worker_id
LEFT JOIN work_order_shares wos ON wos.id = wf.work_order_share_id
`
async function getWorkerOrderFeedbackByIdWithClient(
client: PoolClient,
feedbackId: number | string,
): Promise<WorkerOrderFeedbackRow | null> {
const result = await client.query<WorkerOrderFeedbackRow>(
`${WORKER_ORDER_FEEDBACK_SELECT} WHERE wf.id = $1 LIMIT 1`,
[Number(feedbackId)],
)
return result.rows[0] || null
}
export async function createWorkerOrderFeedback(input: {
workOrderId: number
workerId: number
content: string
now: string
}): Promise<{ feedback: WorkerOrderFeedbackRow | null; failureReason: 'owner_required' | null }> {
return withTransaction(async (client) => {
const orderResult = await client.query<WorkOrderRow>(
'SELECT * FROM work_orders WHERE id = $1 FOR UPDATE',
[input.workOrderId],
)
const order = orderResult.rows[0] || null
if (!order) return { feedback: null, failureReason: 'owner_required' }
const shareResult = await client.query<WorkOrderShareRow>(
`SELECT * FROM work_order_shares
WHERE work_order_id = $1 AND worker_id = $2 AND status != 'cancelled'
FOR UPDATE`,
[input.workOrderId, input.workerId],
)
const share = shareResult.rows[0] || null
const isOrderOwner =
Number(order.assigned_worker_id || 0) === input.workerId ||
Number(order.last_assigned_worker_id || 0) === input.workerId
if (!share && !isOrderOwner) return { feedback: null, failureReason: 'owner_required' }
const result = await client.query<{ id: number }>(
`INSERT INTO worker_order_feedbacks (
work_order_id, work_order_share_id, worker_id, content, status, created_at, updated_at
) VALUES ($1, $2, $3, $4, 'pending', $5, $5)
RETURNING id`,
[input.workOrderId, share?.id || null, input.workerId, input.content, input.now],
)
const feedbackId = Number(result.rows[0]?.id || 0)
await createWorkOrderEventWithClient(client, {
workOrderId: input.workOrderId,
actorType: 'worker',
actorId: String(input.workerId),
eventType: 'worker_feedback_submitted',
fromStatus: order.status,
toStatus: order.status,
payloadJson: JSON.stringify({
feedbackId,
shareId: share?.id || null,
content: input.content,
}),
now: input.now,
})
return {
feedback: await getWorkerOrderFeedbackByIdWithClient(client, feedbackId),
failureReason: null,
}
})
}
export async function listWorkerOrderFeedbacks(
workerId: number | string,
): Promise<WorkerOrderFeedbackRow[]> {
const result = await query<WorkerOrderFeedbackRow>(
`${WORKER_ORDER_FEEDBACK_SELECT}
WHERE wf.worker_id = $1
ORDER BY wf.id DESC
LIMIT 100`,
[Number(workerId)],
)
return result.rows
}
export async function listAdminWorkerOrderFeedbacks(
status = 'pending',
): Promise<WorkerOrderFeedbackRow[]> {
const result = await query<WorkerOrderFeedbackRow>(
`${WORKER_ORDER_FEEDBACK_SELECT}
WHERE ($1 = '' OR wf.status = $1)
ORDER BY wf.id DESC
LIMIT 200`,
[status],
)
return result.rows
}
export async function resolveWorkerOrderFeedback(input: {
feedbackId: number
reply: string
actorName: string
now: string
}): Promise<{
feedback: WorkerOrderFeedbackRow | null
order: WorkOrderRow | null
failureReason: 'not_found' | 'not_pending' | null
}> {
return withTransaction(async (client) => {
const feedbackResult = await client.query<WorkerOrderFeedbackRow>(
'SELECT * FROM worker_order_feedbacks WHERE id = $1 FOR UPDATE',
[input.feedbackId],
)
const feedback = feedbackResult.rows[0] || null
if (!feedback) return { feedback: null, order: null, failureReason: 'not_found' }
if (feedback.status !== 'pending') {
return {
feedback: await getWorkerOrderFeedbackByIdWithClient(client, feedback.id),
order: null,
failureReason: 'not_pending',
}
}
const orderResult = await client.query<WorkOrderRow>(
'SELECT * FROM work_orders WHERE id = $1 FOR UPDATE',
[feedback.work_order_id],
)
const order = orderResult.rows[0] || null
if (!order) return { feedback: null, order: null, failureReason: 'not_found' }
await client.query(
`UPDATE worker_order_feedbacks
SET status = 'resolved', reply = $1, handled_by = $2, handled_at = $3, updated_at = $3
WHERE id = $4`,
[input.reply, input.actorName, input.now, feedback.id],
)
await createWorkOrderEventWithClient(client, {
workOrderId: Number(order.id),
actorType: 'admin',
actorId: input.actorName,
eventType: 'worker_feedback_resolved',
fromStatus: order.status,
toStatus: order.status,
payloadJson: JSON.stringify({ feedbackId: feedback.id, reply: input.reply }),
now: input.now,
})
return {
feedback: await getWorkerOrderFeedbackByIdWithClient(client, feedback.id),
order: await getWorkOrderByIdWithClient(client, order.id),
failureReason: null,
}
})
}
@@ -0,0 +1,82 @@
import { query } from '../../db/client.js'
export async function listWorkerWorkOrderNotes(
workerId: number | string,
workOrderIds: number[],
): Promise<Array<{ work_order_id: number; note: string }>> {
if (workOrderIds.length === 0) return []
const result = await query<{ work_order_id: number; note: string }>(
`
SELECT work_order_id, note
FROM worker_work_order_notes
WHERE worker_id = $1 AND work_order_id = ANY($2::bigint[])
`,
[Number(workerId), workOrderIds],
)
return result.rows
}
export async function upsertWorkerWorkOrderNote(input: {
workOrderId: number
workerId: number
note: string
now: string
}): Promise<string> {
const result = await query<{ note: string }>(
`
INSERT INTO worker_work_order_notes (
work_order_id, worker_id, note, created_at, updated_at
) VALUES ($1, $2, $3, $4, $4)
ON CONFLICT (work_order_id, worker_id) DO UPDATE
SET note = EXCLUDED.note, updated_at = EXCLUDED.updated_at
RETURNING note
`,
[input.workOrderId, input.workerId, input.note, input.now],
)
return String(result.rows[0]?.note || '')
}
/** 查询打手已查看过的工单,用于计算"新"标识。 */
export async function listWorkerWorkOrderViews(
workerId: number | string,
workOrderIds: number[],
): Promise<Array<{ work_order_id: number; viewed_at: string }>> {
if (workOrderIds.length === 0) return []
const result = await query<{ work_order_id: number; viewed_at: string }>(
`
SELECT work_order_id, viewed_at
FROM worker_work_order_views
WHERE worker_id = $1 AND work_order_id = ANY($2::bigint[])
`,
[Number(workerId), workOrderIds],
)
return result.rows
}
/** 首次查看工单时写入已读记录;重复调用保持原查看时间不变。 */
export async function markWorkerWorkOrderViewed(input: {
workOrderId: number
workerId: number
now: string
}): Promise<string> {
const inserted = await query<{ viewed_at: string }>(
`
INSERT INTO worker_work_order_views (work_order_id, worker_id, viewed_at)
VALUES ($1, $2, $3)
ON CONFLICT (work_order_id, worker_id) DO NOTHING
RETURNING viewed_at
`,
[input.workOrderId, input.workerId, input.now],
)
if (inserted.rows[0]?.viewed_at) return String(inserted.rows[0].viewed_at)
const existing = await query<{ viewed_at: string }>(
`
SELECT viewed_at
FROM worker_work_order_views
WHERE work_order_id = $1 AND worker_id = $2
LIMIT 1
`,
[input.workOrderId, input.workerId],
)
return String(existing.rows[0]?.viewed_at || input.now)
}
@@ -0,0 +1,455 @@
import { listKuaishouIndustryVoucherWorkOrderSyncSources } from '../../repositories/kuaishou-industry-voucher-repo.js'
import {
deleteWorkProductRuleMapping,
listWorkProductMatchLogs,
listWorkProductRuleMappings,
listWorkProductRuleMappingsPage,
upsertWorkProductRuleMapping,
} from '../../repositories/worker-platform/product-match-repo.js'
import { listWorkProductRules } from '../../repositories/worker-platform/work-product-rule-repo.js'
import type {
WorkProductMatchLogRow,
WorkProductRuleMappingRow,
} from '../../repositories/worker-platform/types.js'
import type { JsonObject } from '../../types/json.js'
import type { OrderItemRow, OrderRow } from '../../types/repository/rows.js'
import { createHttpError } from '../../utils/http.js'
import { nowIso } from '../../utils/time.js'
import { normalizePage, normalizePageSize, safeParseJson } from '../admin/admin-query-utils.js'
import {
mapWorkProductRule,
normalizeBoolean,
normalizeOptionalId,
normalizePositiveInteger,
resolveKuaishouWorkProductMatchContext,
resolveMatchingProductRuleDecision,
} from './mappers.js'
import { reprocessKuaishouSendCodeWorkOrders } from './sync-work-orders-from-send-code.js'
import { getWorkerProductMatchConfig } from './worker-product-match-config-service.js'
export async function reprocessAdminKuaishouSendCodeWorkOrders(payload: JsonObject = {}) {
const limit = Math.min(500, normalizePositiveInteger(payload.limit, 100))
return reprocessKuaishouSendCodeWorkOrders(limit)
}
export async function listAdminWorkProductRuleMappings(query: JsonObject = {}) {
const ruleId = normalizeOptionalId(query.ruleId ?? query.rule_id) || 0
const hasPagination = query.page !== undefined || query.pageSize !== undefined || query.keyword
if (hasPagination) {
const page = normalizePage(query.page)
const pageSize = normalizePageSize(query.pageSize)
const keyword = String(query.keyword || '').trim()
const result = await listWorkProductRuleMappingsPage({
page,
pageSize,
keyword,
ruleId,
})
return {
items: result.items.map(mapWorkProductRuleMapping),
pagination: { page, pageSize, total: result.total },
}
}
const items = await listWorkProductRuleMappings({ ruleId })
return { items: items.map(mapWorkProductRuleMapping) }
}
export async function saveAdminWorkProductRuleMapping(payload: JsonObject = {}) {
const ruleId = normalizeOptionalId(payload.ruleId ?? payload.rule_id)
const sellerIds = normalizeKuaishouMatchIds(
payload.sellerIds ?? payload.seller_ids ?? payload.sellerId ?? payload.seller_id,
)
const relItemIds = normalizeKuaishouMatchIdList(payload.relItemId ?? payload.rel_item_id)
const relItemId = relItemIds.join(',')
const itemTitle = String(payload.itemTitle ?? payload.item_title ?? '').trim()
const rawMappingType = String(payload.mappingType ?? payload.mapping_type ?? '').trim()
const mappingType =
rawMappingType === 'sku_series'
? 'sku_series'
: rawMappingType === 'sku_exact' || rawMappingType === 'sku_override'
? 'sku_exact'
: 'product_default'
const relSkuId =
mappingType === 'sku_exact'
? normalizeKuaishouMatchId(payload.relSkuId ?? payload.rel_sku_id)
: ''
const skuNick =
mappingType === 'product_default'
? ''
: String(payload.skuNick ?? payload.sku_nick ?? '').trim()
const skuMatchMode =
mappingType === 'sku_series' &&
String(payload.skuMatchMode ?? payload.sku_match_mode ?? '').trim() === 'exact'
? 'exact'
: 'fuzzy'
const requiresProductScope = mappingType === 'product_default'
if (
!ruleId ||
sellerIds.length === 0 ||
(requiresProductScope && relItemIds.length === 0 && !itemTitle)
) {
throw createHttpError(
'请选择接单模板,并填写至少一个店铺;商品默认规则还需填写大标题或关联商品 ID',
{
statusCode: 400,
errorCode: 'work_product_mapping_identity_required',
},
)
}
if (
mappingType === 'product_default' &&
sellerIds.length > 1 &&
(!itemTitle || relItemIds.length > 0)
) {
throw createHttpError('多店共享规则请填写快手大标题,并留空关联商品 ID', {
statusCode: 400,
errorCode: 'work_product_mapping_multi_shop_scope_invalid',
})
}
if (mappingType === 'sku_exact' && !relSkuId && !skuNick) {
throw createHttpError('SKU 精确规则需要关联 SKU ID 或具体 SKU 名称', {
statusCode: 400,
errorCode: 'work_product_mapping_sku_required',
})
}
if (mappingType === 'sku_series' && !skuNick) {
throw createHttpError('SKU 数量系列需要填写系列名称', {
statusCode: 400,
errorCode: 'work_product_mapping_series_required',
})
}
const rule = (await listWorkProductRules()).find((item) => Number(item.id) === ruleId)
if (!rule) {
throw createHttpError('接单模板不存在', {
statusCode: 404,
errorCode: 'work_product_mapping_rule_not_found',
})
}
const mappingId = normalizeOptionalId(payload.mappingId ?? payload.mapping_id)
if (normalizeBoolean(payload.enabled, true)) {
await assertNoConflictingWorkProductRuleMapping({
mappingId,
sellerIds,
relItemId,
itemTitle,
relSkuId,
skuNick,
mappingType,
})
}
const mapping = await upsertWorkProductRuleMapping({
mappingId,
ruleId,
sellerIds,
relItemId,
itemTitle,
relSkuId,
skuNick,
skuMatchMode,
mappingType,
enabled: normalizeBoolean(payload.enabled, true),
now: nowIso(),
})
if (!mapping) {
throw createHttpError('商品映射保存失败', {
statusCode: 500,
errorCode: 'work_product_mapping_save_failed',
})
}
return {
mapping: mapWorkProductRuleMapping({
...mapping,
rule_key: rule.rule_key,
product_name: rule.product_name,
}),
}
}
export async function deleteAdminWorkProductRuleMapping(mappingId: number | string) {
const normalizedMappingId = normalizeOptionalId(mappingId)
if (!normalizedMappingId) {
throw createHttpError('商品映射 ID 不正确', {
statusCode: 400,
errorCode: 'work_product_mapping_id_invalid',
})
}
return deleteWorkProductRuleMapping(normalizedMappingId)
}
export async function listAdminKuaishouMatchSources(payload: JsonObject = {}) {
const limit = Math.min(500, normalizePositiveInteger(payload.limit, 100))
const sources = await listKuaishouIndustryVoucherWorkOrderSyncSources(limit)
const grouped = new Map<string, JsonObject>()
for (const source of sources) {
const rawPayload = safeParseJson(source.raw_payload_json)
const body = resolveKuaishouMatchPayload(rawPayload)
const ext = safeParseJson(body.ext)
const sellerId = normalizeKuaishouMatchId(body.sellerId)
const relItemId = normalizeKuaishouMatchId(ext.relItemId)
const itemTitle = String(body.itemTitle || '').trim()
const relSkuId = normalizeKuaishouMatchId(ext.relSkuId)
const skuNick = String(ext.skuNick || '').trim()
if (!sellerId || (!relItemId && !itemTitle)) continue
const key = [
sellerId,
relItemId,
normalizeMatchTextForAdmin(itemTitle),
relSkuId,
normalizeMatchTextForAdmin(skuNick),
].join('|')
const current = grouped.get(key)
grouped.set(key, {
sellerId,
relItemId,
itemTitle,
relSkuId,
skuNick,
sampleOid: String(body.oid || source.oid || '').trim(),
lastSeenAt: source.updated_at,
seenCount: Number(current?.seenCount || 0) + 1,
rawPayload,
})
}
return {
items: Array.from(grouped.values()).sort((left, right) =>
String(right.lastSeenAt).localeCompare(String(left.lastSeenAt)),
),
}
}
export async function testAdminKuaishouProductMatch(payload: JsonObject = {}) {
const body = resolveKuaishouMatchPayload(payload.rawPayload ?? payload.raw_payload ?? payload)
const ext = safeParseJson(body.ext)
const order = {
provider: '91kaquan',
platform: 'kuaishou',
shop_id: String(body.sellerId || '').trim(),
} as OrderRow
const item = {
id: 0,
order_id: 0,
sku_code: String(body.itemId || '').trim(),
sku_name: String(ext.skuNick || '').trim(),
quantity: Math.max(1, Number(body.num) || 1),
spec_json: {},
item_snapshot_json: {
kuaishouSendCode: {
sellerId: String(body.sellerId || '').trim(),
itemId: String(body.itemId || '').trim(),
itemTitle: String(body.itemTitle || '').trim(),
skuId: String(body.skuId || '').trim(),
relItemId: normalizeKuaishouMatchId(ext.relItemId),
relSkuId: normalizeKuaishouMatchId(ext.relSkuId),
skuNick: String(ext.skuNick || '').trim(),
},
},
} as OrderItemRow
const [rules, mappings] = await Promise.all([
listWorkProductRules({ enabled: true }),
listWorkProductRuleMappings({ enabled: true }),
])
const decision = resolveMatchingProductRuleDecision(order, item, rules, mappings, {
legacyFallbackEnabled: getWorkerProductMatchConfig().legacyFallbackEnabled,
})
const context = resolveKuaishouWorkProductMatchContext(item)
return {
context,
status: decision.reason,
mappingId: decision.mappingId,
rule: decision.rule ? mapWorkProductRule(decision.rule) : null,
candidates: decision.candidates,
}
}
export async function listAdminWorkProductMatchLogs(payload: JsonObject = {}) {
const limit = Math.min(500, normalizePositiveInteger(payload.limit, 100))
return { items: (await listWorkProductMatchLogs(limit)).map(mapWorkProductMatchLog) }
}
function mapWorkProductRuleMapping(mapping: WorkProductRuleMappingRow) {
const sellerIds = parseKuaishouMappingSellerIds(mapping.seller_ids_json)
return {
mappingId: Number(mapping.id),
ruleId: Number(mapping.rule_id),
ruleKey: String(mapping.rule_key || ''),
productName: String(mapping.product_name || ''),
sellerIds,
sellerId: sellerIds[0] || '',
relItemId: String(mapping.rel_item_id || ''),
itemTitle: String(mapping.item_title || ''),
relSkuId: String(mapping.rel_sku_id || ''),
skuNick: String(mapping.sku_nick || ''),
skuMatchMode: mapping.sku_match_mode === 'exact' ? 'exact' : 'fuzzy',
mappingType:
mapping.mapping_type === 'sku_series'
? 'sku_series'
: mapping.mapping_type === 'sku_exact' || mapping.mapping_type === 'sku_override'
? 'sku_exact'
: 'product_default',
enabled: mapping.enabled === true,
createdAt: mapping.created_at,
updatedAt: mapping.updated_at,
}
}
function mapWorkProductMatchLog(log: WorkProductMatchLogRow) {
return {
logId: Number(log.id),
orderId: log.order_id ? Number(log.order_id) : null,
orderItemId: log.order_item_id ? Number(log.order_item_id) : null,
source: String(log.source || ''),
sellerId: String(log.seller_id || ''),
relItemId: String(log.rel_item_id || ''),
itemTitle: String(log.item_title || ''),
relSkuId: String(log.rel_sku_id || ''),
skuNick: String(log.sku_nick || ''),
status: String(log.match_status || ''),
ruleId: log.rule_id ? Number(log.rule_id) : null,
mappingId: log.mapping_id ? Number(log.mapping_id) : null,
ruleKey: String(log.rule_key || ''),
productName: String(log.product_name || ''),
candidates: parseJsonArray(log.candidates_json),
rawPayload: safeParseJson(log.raw_payload_json),
createdAt: log.created_at,
}
}
function resolveKuaishouMatchPayload(value: unknown): JsonObject {
const raw = safeParseJson(value)
const body = safeParseJson(raw.body)
return Object.keys(body).length > 0 ? body : raw
}
function normalizeKuaishouMatchId(value: unknown) {
const normalized = String(value || '').trim()
return normalized === '0' ? '' : normalized
}
function normalizeKuaishouMatchIds(value: unknown) {
const values = Array.isArray(value)
? value
: String(value ?? '')
.split(/[,\n]/)
.map((item) => item.trim())
return Array.from(new Set(values.map(normalizeKuaishouMatchId).filter(Boolean)))
}
function normalizeKuaishouMatchIdList(value: unknown) {
const values = Array.isArray(value) ? value : [value]
return normalizeKuaishouMatchIds(
values.flatMap((item) => String(item ?? '').split(/[,;\s\n\r]+/)),
)
}
function parseKuaishouMappingSellerIds(value: unknown) {
if (Array.isArray(value)) return normalizeKuaishouMatchIds(value)
try {
return normalizeKuaishouMatchIds(JSON.parse(String(value || '[]')))
} catch {
return normalizeKuaishouMatchIds(value)
}
}
function normalizeMatchTextForAdmin(value: unknown) {
return String(value || '')
.normalize('NFKC')
.toLowerCase()
.replace(/[\s\u3000]+/g, '')
.trim()
}
function assertNoConflictingWorkProductRuleMapping(input: {
mappingId: number | null
sellerIds: string[]
relItemId: string
itemTitle: string
relSkuId: string
skuNick: string
mappingType: 'product_default' | 'sku_exact' | 'sku_series'
}) {
return listWorkProductRuleMappings({ enabled: true }).then((mappings) => {
const conflict = mappings.find((mapping) => {
if (input.mappingId && Number(mapping.id) === input.mappingId) return false
const existingType =
mapping.mapping_type === 'sku_series'
? 'sku_series'
: mapping.mapping_type === 'sku_exact' || mapping.mapping_type === 'sku_override'
? 'sku_exact'
: 'product_default'
if (existingType !== input.mappingType) return false
const existingSellerIds = parseKuaishouMappingSellerIds(mapping.seller_ids_json)
if (!existingSellerIds.some((sellerId) => input.sellerIds.includes(sellerId))) return false
const inputItemIds = normalizeKuaishouMatchIdList(input.relItemId)
const existingItemIds = normalizeKuaishouMatchIdList(mapping.rel_item_id)
if (
!mappingProductScopesOverlap(
inputItemIds,
input.itemTitle,
existingItemIds,
mapping.item_title,
)
) {
return false
}
if (input.mappingType === 'product_default') return true
if (input.mappingType === 'sku_series') {
// 精确系列是模糊系列的子集,两者并存仍可能产生同分冲突,因此统一拦截。
return (
normalizeMatchTextForAdmin(input.skuNick) === normalizeMatchTextForAdmin(mapping.sku_nick)
)
}
const inputSkuId = normalizeKuaishouMatchId(input.relSkuId)
const existingSkuId = normalizeKuaishouMatchId(mapping.rel_sku_id)
const skuIdConflict = Boolean(inputSkuId && existingSkuId && inputSkuId === existingSkuId)
const skuNameConflict =
Boolean(input.skuNick && mapping.sku_nick) &&
normalizeMatchTextForAdmin(input.skuNick) === normalizeMatchTextForAdmin(mapping.sku_nick)
return skuIdConflict || skuNameConflict
})
if (conflict) {
throw createHttpError('商品映射与已有启用映射冲突,请停用或调整其中一条', {
statusCode: 409,
errorCode: 'work_product_mapping_conflict',
context: {
mappingId: Number(conflict.id),
ruleKey: String(conflict.rule_key || ''),
},
})
}
})
}
function mappingProductScopesOverlap(
leftItemIds: string[],
leftTitle: string,
rightItemIds: string[],
rightTitle: string,
) {
if (leftItemIds.length > 0 && rightItemIds.length > 0) {
return leftItemIds.some((itemId) => rightItemIds.includes(itemId))
}
// 只有一侧指定商品 ID 时,另一侧可能是全商品范围,按可能重叠处理并拦截。
if (leftItemIds.length > 0 || rightItemIds.length > 0) return true
const normalizedLeftTitle = normalizeMatchTextForAdmin(leftTitle)
const normalizedRightTitle = normalizeMatchTextForAdmin(rightTitle)
if (normalizedLeftTitle || normalizedRightTitle) {
return Boolean(
normalizedLeftTitle && normalizedRightTitle && normalizedLeftTitle === normalizedRightTitle,
)
}
return true
}
function parseJsonArray(value: unknown): unknown[] {
if (Array.isArray(value)) return value
try {
const parsed = JSON.parse(String(value || '[]'))
return Array.isArray(parsed) ? parsed : []
} catch {
return []
}
}
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,474 @@
import {
countWorkCategoryUsages,
deleteWorkCategory,
getWorkCategoryById,
listAllWorkCategories,
upsertWorkCategory,
} from '../../repositories/worker-platform/work-category-repo.js'
import {
countWorkerLevelUsages,
deleteWorkerLevel,
getWorkerLevelById,
getWorkerLevelByKey,
listWorkerLevels,
upsertWorkerLevel,
} from '../../repositories/worker-platform/worker-repo.js'
import {
countWorkProductRulesByCategory,
deleteWorkProductRule,
getWorkProductRuleById,
listOpenWorkOrdersByProductRuleId,
listWorkProductRulesPage,
syncOpenWorkOrderPriceFromRule,
upsertWorkProductRule,
} from '../../repositories/worker-platform/work-product-rule-repo.js'
import type { WorkOrderRow, WorkProductRuleRow } from '../../repositories/worker-platform/types.js'
import type { JsonObject } from '../../types/json.js'
import { createHttpError } from '../../utils/http.js'
import { randomId } from '../../utils/random.js'
import { nowIso } from '../../utils/time.js'
import { normalizePage, normalizePageSize, safeParseJson } from '../admin/admin-query-utils.js'
import { publishWorkOrderRealtimeChange } from '../realtime/realtime-event-service.js'
import {
DEFAULT_CATEGORY_KEY,
DEFAULT_LEVEL_KEY,
DEFAULT_LEVEL_NAME,
mapWorkCategory,
mapWorkProductRule,
mapWorkerLevel,
normalizeAmountFen,
normalizeBoolean,
normalizeEnabledStatus,
normalizeInteger,
normalizeMatchType,
normalizeOptionalId,
normalizePositiveInteger,
normalizeRequirementFieldsFromPayload,
normalizeSlugKey,
normalizeWorkProductRuleMatch,
resolveSkuNameQuantity,
resolveWorkOrderRulePricing,
} from './mappers.js'
import { ensureWorkerPlatformDefaults, normalizeWorkOrderTimeoutPolicy } from './worker-service.js'
export async function listAdminWorkerLevels() {
await ensureWorkerPlatformDefaults()
return { items: (await listWorkerLevels()).map(mapWorkerLevel) }
}
export async function listAdminWorkCategories() {
await ensureWorkerPlatformDefaults()
return { items: (await listAllWorkCategories()).map(mapWorkCategory) }
}
export async function saveAdminWorkCategory(payload: JsonObject = {}) {
const now = nowIso()
const categoryKey = normalizeSlugKey(
payload.categoryKey || payload.category_key || payload.name,
DEFAULT_CATEGORY_KEY,
)
const name = String(payload.name || '').trim() || '默认分类'
const category = await upsertWorkCategory({
categoryKey,
name,
sortOrder: normalizeInteger(payload.sortOrder, 100),
status: normalizeEnabledStatus(payload.status),
now,
})
return { category: mapWorkCategory(category) }
}
export async function deleteAdminWorkCategory(categoryId: number | string) {
const category = await getWorkCategoryById(Number(categoryId))
if (!category) {
throw createHttpError('分类不存在', {
statusCode: 404,
errorCode: 'work_category_not_found',
})
}
const usages = await countWorkCategoryUsages(Number(categoryId))
if (usages.workOrderCount > 0 || usages.productRuleCount > 0) {
throw createHttpError(
`该分类已被 ${usages.workOrderCount} 个工单和 ${usages.productRuleCount} 条物品规则使用,无法删除,请先停用`,
{
statusCode: 409,
errorCode: 'work_category_in_use',
},
)
}
const { deleted } = await deleteWorkCategory(Number(categoryId))
return { deleted }
}
export async function saveAdminWorkerLevel(payload: JsonObject = {}) {
const now = nowIso()
const levelKey = String(payload.levelKey || payload.level_key || '').trim() || DEFAULT_LEVEL_KEY
const existingLevel = await getWorkerLevelByKey(levelKey)
const existingPermissions = safeParseJson(existingLevel?.permission_json)
const name = String(payload.name || '').trim() || DEFAULT_LEVEL_NAME
const permissions = {
depositFreeAmount: normalizeAmountFen(
payload.depositFreeAmount ?? payload.depositFreeAmountYuan,
0,
),
maxActiveOrders: normalizePositiveInteger(payload.maxActiveOrders, 1),
upgradeThreshold: normalizeInteger(payload.upgradeThreshold, 0),
autoAcceptWithoutEvidence: normalizeBoolean(
payload.autoAcceptWithoutEvidence,
normalizeBoolean(existingPermissions.autoAcceptWithoutEvidence, levelKey === 'vip5'),
),
visibleDelaySeconds: normalizeInteger(payload.visibleDelaySeconds, -1),
}
const level = await upsertWorkerLevel({
levelKey,
name,
sortOrder: normalizeInteger(payload.sortOrder, Number(existingLevel?.sort_order || 100)),
status: String(payload.status || 'active').trim() === 'disabled' ? 'disabled' : 'active',
permissionJson: JSON.stringify(permissions),
now,
})
return { level: mapWorkerLevel(level) }
}
export async function deleteAdminWorkerLevel(levelId: number | string) {
const level = await getWorkerLevelById(Number(levelId))
if (!level) {
throw createHttpError('等级不存在', {
statusCode: 404,
errorCode: 'worker_level_not_found',
})
}
if (level.level_key === DEFAULT_LEVEL_KEY) {
throw createHttpError('默认等级(VIP1)不可删除', {
statusCode: 409,
errorCode: 'worker_level_default_protected',
})
}
const usageCount = await countWorkerLevelUsages(Number(levelId))
if (usageCount > 0) {
throw createHttpError(`该等级正在被 ${usageCount} 名打手使用,请先调整打手等级后再删除`, {
statusCode: 409,
errorCode: 'worker_level_in_use',
})
}
const { deleted } = await deleteWorkerLevel(Number(levelId))
return { deleted }
}
export async function listAdminWorkProductRules(query: JsonObject = {}) {
await ensureWorkerPlatformDefaults()
const page = normalizePage(query.page)
const pageSize = normalizePageSize(query.pageSize)
const enabledValue = String(query.enabled ?? '').trim()
const enabled = enabledValue
? ['true', '1', 'enabled', 'active'].includes(enabledValue.toLowerCase())
: null
const keyword = String(query.keyword || '').trim()
const categoryId = normalizeOptionalId(query.categoryId ?? query.category_id) || 0
const [{ items, total }, categoryRows] = await Promise.all([
listWorkProductRulesPage({
page,
pageSize,
enabled,
keyword,
categoryId,
}),
countWorkProductRulesByCategory({ enabled, keyword }),
])
return {
items: items.map(mapWorkProductRule),
pagination: { page, pageSize, total },
categoryCounts: categoryRows.map((row) => ({
categoryId: row.categoryId,
total: row.total,
})),
}
}
export function resolveSharingQuantity(
totalAmount: number,
fallbackQuantity: number,
unitReward: number,
): number {
if (totalAmount > 0 && unitReward > 0) {
return Math.max(1, Math.round(totalAmount / unitReward))
}
return Math.max(1, fallbackQuantity)
}
function hasConfiguredSharingValue(value: unknown) {
return value !== undefined && value !== null && String(value).trim() !== ''
}
function hasOwnPayload(payload: JsonObject, key: string) {
return Object.prototype.hasOwnProperty.call(payload, key)
}
export async function saveAdminWorkProductRule(payload: JsonObject = {}) {
const defaults = await ensureWorkerPlatformDefaults()
const ruleId = normalizeOptionalId(payload.ruleId ?? payload.rule_id)
const existingRule = ruleId ? await getWorkProductRuleById(ruleId) : null
if (ruleId && !existingRule) {
throw createHttpError('接单模板不存在', {
statusCode: 404,
errorCode: 'work_product_rule_not_found',
})
}
const hasLegacyMatchPayload = [
'match',
'sellerId',
'seller_id',
'itemId',
'item_id',
'relItemId',
'rel_item_id',
'skuId',
'sku_id',
'relSkuId',
'rel_sku_id',
'itemTitle',
'item_title',
'skuNick',
'sku_nick',
'itemTitleMatchType',
'item_title_match_type',
'skuNickMatchType',
'sku_nick_match_type',
].some((key) => Object.prototype.hasOwnProperty.call(payload, key))
const match = hasLegacyMatchPayload
? normalizeWorkProductRuleMatch(payload)
: safeParseJson(existingRule?.match_json)
const productName = String(
payload.productName || payload.product_name || payload.skuName || '',
).trim()
const skuCode = String(payload.skuCode || payload.sku_code || '').trim()
const matchedProductName = String(match.itemTitle || match.skuNick || '').trim()
if (!productName && !skuCode && !matchedProductName) {
throw createHttpError('请填写商品名或 SKU', {
statusCode: 400,
errorCode: 'work_product_rule_target_required',
})
}
const rewardAmount = normalizeAmountFen(payload.rewardAmount ?? payload.rewardAmountYuan, 0)
const unitPriceFen =
payload.unitPrice === undefined && payload.unitPriceYuan === undefined
? 0
: normalizeAmountFen(payload.unitPrice ?? payload.unitPriceYuan, 0)
const sharingEnabled = normalizeBoolean(payload.sharingEnabled ?? payload.sharing_enabled, false)
const hasManualSharingConfig =
hasConfiguredSharingValue(payload.sharingTotalQuantity ?? payload.sharing_total_quantity) ||
hasConfiguredSharingValue(payload.sharingUnitReward ?? payload.sharingUnitRewardYuan) ||
hasConfiguredSharingValue(payload.sharingTotalAmount ?? payload.sharingTotalAmountYuan)
const sharingAutoFromOrder =
sharingEnabled &&
unitPriceFen > 0 &&
(normalizeBoolean(payload.sharingAutoFromOrder ?? payload.sharing_auto_from_order, false) ||
!hasManualSharingConfig)
const hasManualSharingQuantity = hasConfiguredSharingValue(
payload.sharingTotalQuantity ?? payload.sharing_total_quantity,
)
const hasManualSharingUnitReward = hasConfiguredSharingValue(
payload.sharingUnitReward ?? payload.sharingUnitRewardYuan,
)
if (
sharingEnabled &&
unitPriceFen > 0 &&
!sharingAutoFromOrder &&
hasManualSharingQuantity !== hasManualSharingUnitReward
) {
throw createHttpError('按件计价的拼单覆盖值请同时填写拼单数量和拼单单价', {
statusCode: 400,
errorCode: 'work_product_rule_sharing_override_incomplete',
})
}
const sharingUnitReward = normalizeAmountFen(
payload.sharingUnitReward ?? payload.sharingUnitRewardYuan,
0,
)
const sharingTotalAmount = normalizeAmountFen(
payload.sharingTotalAmount ?? payload.sharingTotalAmountYuan,
0,
)
const sharingTotalQuantity = resolveSharingQuantity(
sharingTotalAmount,
normalizePositiveInteger(payload.sharingTotalQuantity ?? payload.sharing_total_quantity, 1),
sharingUnitReward,
)
const resolvedSharingTotalAmount =
sharingTotalAmount > 0 ? sharingTotalAmount : sharingUnitReward * sharingTotalQuantity
if (sharingEnabled && !sharingAutoFromOrder) {
if (resolvedSharingTotalAmount <= 0) {
throw createHttpError('启用拼单时请填写拼单总价或单价', {
statusCode: 400,
errorCode: 'work_product_rule_sharing_amount_invalid',
})
}
}
const finalRewardAmount = sharingEnabled
? sharingAutoFromOrder
? 0
: resolvedSharingTotalAmount
: rewardAmount
const fields = normalizeRequirementFieldsFromPayload(payload)
const ruleKey = existingRule?.rule_key || randomId('rule-')
const rule = await upsertWorkProductRule({
ruleKey,
provider: hasOwnPayload(payload, 'provider')
? String(payload.provider || '').trim()
: String(existingRule?.provider || '').trim(),
platform: hasOwnPayload(payload, 'platform')
? String(payload.platform || '').trim()
: String(existingRule?.platform || '').trim(),
shopId:
hasOwnPayload(payload, 'shopId') || hasOwnPayload(payload, 'shop_id')
? String(payload.shopId || payload.shop_id || '').trim()
: String(existingRule?.shop_id || '').trim(),
skuCode:
hasOwnPayload(payload, 'skuCode') || hasOwnPayload(payload, 'sku_code')
? skuCode
: String(existingRule?.sku_code || '').trim(),
productName: productName || matchedProductName,
matchType:
hasOwnPayload(payload, 'matchType') || hasOwnPayload(payload, 'match_type')
? normalizeMatchType(payload.matchType || payload.match_type)
: normalizeMatchType(existingRule?.match_type),
categoryId:
normalizeOptionalId(payload.categoryId || payload.category_id) ||
defaults.category?.id ||
null,
enabled: normalizeBoolean(payload.enabled, true),
autoCreate: normalizeBoolean(payload.autoCreate ?? payload.auto_create, true),
rewardAmount: finalRewardAmount,
unitPriceFen,
// 商品规则不配置押金;工单创建后由运营按需手工填写。
requiredDepositAmount: 0,
depositThresholdAmount: 0,
sharingEnabled,
sharingTotalQuantity,
sharingUnitReward,
sharingAutoFromOrder,
timeoutMinutes: Math.max(0, normalizeInteger(payload.timeoutMinutes, 0)),
timeoutPolicy: normalizeWorkOrderTimeoutPolicy(payload.timeoutPolicy),
requirementJson: JSON.stringify({ fields }),
matchJson: JSON.stringify(match),
sortOrder: normalizeInteger(payload.sortOrder, 100),
now: nowIso(),
})
return { rule: mapWorkProductRule(rule) }
}
/** 根据工单创建时的规格,计算当前模板可同步的价格。 */
export function resolveWorkOrderTemplateSyncPrice(input: {
rule: WorkProductRuleRow
workOrder: WorkOrderRow
}): { rewardAmount: number; sharingUnitReward?: number } | null {
const material = safeParseJson(input.workOrder.material_json)
const source = safeParseJson(material.source)
const sourceSkuQuantity = Number(source.skuQuantity || 0)
const skuQuantity =
Number.isSafeInteger(sourceSkuQuantity) && sourceSkuQuantity > 0
? sourceSkuQuantity
: resolveSkuNameQuantity(input.workOrder.product_name)
const pricing = resolveWorkOrderRulePricing({
skuQuantity,
unitPriceFen: Number(input.rule.unit_price_fen || 0),
fixedRewardAmount: Number(input.rule.reward_amount || 0),
sharingEnabled: input.rule.sharing_enabled === true,
sharingAutoFromOrder: input.rule.sharing_auto_from_order === true,
sharingTotalQuantity: Number(input.rule.sharing_total_quantity || 1),
sharingUnitReward: Number(input.rule.sharing_unit_reward || 0),
})
if (input.workOrder.sharing_enabled === true) {
if (input.rule.sharing_enabled !== true || pricing.sharingUnitReward <= 0) {
return null
}
return {
rewardAmount: pricing.rewardAmount,
sharingUnitReward: pricing.sharingUnitReward,
}
}
return pricing.rewardAmount > 0 ? { rewardAmount: pricing.rewardAmount } : null
}
/** 将指定模板的当前价格同步到仍在接单大厅的来源工单。 */
export async function syncAdminWorkProductRulePrice(ruleId: number | string, actorName = '') {
const normalizedRuleId = normalizeOptionalId(ruleId)
const rule = normalizedRuleId ? await getWorkProductRuleById(normalizedRuleId) : null
if (!rule) {
throw createHttpError('接单模板不存在', {
statusCode: 404,
errorCode: 'work_product_rule_not_found',
})
}
const candidates = await listOpenWorkOrdersByProductRuleId(rule.id)
const result = {
candidateCount: candidates.length,
updatedCount: 0,
updatedNormalCount: 0,
updatedPartialSharingCount: 0,
skippedInvalidPricingCount: 0,
skippedUnchangedCount: 0,
skippedUnavailableCount: 0,
}
const now = nowIso()
for (const workOrder of candidates) {
const price = resolveWorkOrderTemplateSyncPrice({ rule, workOrder })
if (!price) {
result.skippedInvalidPricingCount += 1
continue
}
const synced = await syncOpenWorkOrderPriceFromRule({
workOrderId: Number(workOrder.id),
productRuleId: Number(rule.id),
rewardAmount: price.rewardAmount,
...(price.sharingUnitReward === undefined
? {}
: { sharingUnitReward: price.sharingUnitReward }),
actorName,
now,
})
if (!synced.updated) {
if (synced.reason === 'price_unchanged') {
result.skippedUnchangedCount += 1
} else {
result.skippedUnavailableCount += 1
}
continue
}
result.updatedCount += 1
if (synced.sharingPartial) {
result.updatedPartialSharingCount += 1
} else {
result.updatedNormalCount += 1
}
publishWorkOrderRealtimeChange({
workOrderId: synced.workOrderId,
workerIds: synced.workerIds,
hallChanged: true,
})
}
return { rule: mapWorkProductRule(rule), ...result }
}
export async function deleteAdminWorkProductRule(ruleId: number | string) {
const normalizedRuleId = normalizeOptionalId(ruleId)
if (!normalizedRuleId) {
throw createHttpError('物品规则 ID 不正确', {
statusCode: 400,
errorCode: 'work_product_rule_id_invalid',
})
}
const { deleted } = await deleteWorkProductRule(normalizedRuleId)
if (!deleted) {
throw createHttpError('物品规则不存在或已删除', {
statusCode: 404,
errorCode: 'work_product_rule_not_found',
})
}
return { deleted }
}
@@ -0,0 +1,37 @@
import type { JsonObject } from '../../types/json.js'
import {
getWorkerAnnouncementConfig,
saveWorkerAnnouncementConfig,
} from './worker-announcement-config-service.js'
import {
getWorkerPlatformNotificationConfig,
saveWorkerPlatformNotificationConfig,
} from './worker-platform-notification-config-service.js'
import {
getWorkerProductMatchConfig,
saveWorkerProductMatchConfig,
} from './worker-product-match-config-service.js'
export async function getAdminWorkerPlatformNotificationConfig() {
return getWorkerPlatformNotificationConfig()
}
export async function saveAdminWorkerPlatformNotificationConfig(payload: JsonObject = {}) {
return saveWorkerPlatformNotificationConfig(payload)
}
export async function getAdminWorkerProductMatchConfig() {
return getWorkerProductMatchConfig()
}
export async function getAdminWorkerAnnouncementConfig() {
return getWorkerAnnouncementConfig()
}
export async function saveAdminWorkerAnnouncementConfig(payload: JsonObject = {}) {
return saveWorkerAnnouncementConfig(payload)
}
export async function saveAdminWorkerProductMatchConfig(payload: JsonObject = {}) {
return saveWorkerProductMatchConfig(payload)
}
@@ -0,0 +1,235 @@
import {
addWorkerWalletCredit,
getWorkerFinanceRequestById,
listWorkerFinanceRequests,
listWorkerWithdrawalAccounts,
reviewWorkerFinanceRequest,
upsertWorkerWithdrawalAccount,
} from '../../repositories/worker-platform/worker-repo.js'
import type { WorkerWithdrawalAccountRow } from '../../repositories/worker-platform/types.js'
import type { JsonObject } from '../../types/json.js'
import { createHttpError } from '../../utils/http.js'
import { nowIso } from '../../utils/time.js'
import { normalizePage, normalizePageSize, safeParseJson } from '../admin/admin-query-utils.js'
import { resolveAdminNotificationEntity } from '../admin/admin-notification-service.js'
import { refreshUploadedFileUrls } from '../file-storage/file-storage-service.js'
import {
publishWorkerFinanceRealtimeChange,
publishWorkerWalletRealtimeChange,
} from '../realtime/realtime-event-service.js'
import {
mapFinanceRequest,
mapWallet,
mapWorkerUser,
normalizeAdminFinanceReviewStatus,
normalizeAmountFen,
normalizeFinanceRequestStatus,
normalizeFinanceRequestType,
normalizeOptionalId,
normalizeProofFiles,
normalizeWithdrawChannel,
} from './mappers.js'
import { getRequiredWorker } from './worker-service.js'
import { getWorkerFinanceConfig, saveWorkerFinanceConfig } from './worker-finance-config-service.js'
function mapAdminWorkerWithdrawalAccount(account: WorkerWithdrawalAccountRow) {
const alipayQrCodeImage = refreshUploadedFileUrls(safeParseJson(account.alipay_qr_code_json))
const wechatQrCodeImage = refreshUploadedFileUrls(safeParseJson(account.wechat_qr_code_json))
return {
accountChannel: account.account_channel,
accountName: account.account_name,
accountNo: account.account_channel === 'alipay' ? account.account_no : '',
alipayQrCodeImage: String(alipayQrCodeImage.url || '').trim() ? alipayQrCodeImage : null,
wechatQrCodeImage: String(wechatQrCodeImage.url || '').trim() ? wechatQrCodeImage : null,
createdAt: account.created_at,
}
}
export async function creditAdminWorkerWallet(workerId: number | string, payload: JsonObject = {}) {
await getRequiredWorker(workerId)
const amount = normalizeAmountFen(payload.amount ?? payload.amountYuan, 0)
if (amount <= 0) {
throw createHttpError('充值金额必须大于 0', {
statusCode: 400,
errorCode: 'worker_credit_amount_invalid',
})
}
const wallet = await addWorkerWalletCredit({
workerId: Number(workerId),
amount,
note: String(payload.note || '后台人工充值').trim(),
payloadJson: JSON.stringify({ source: 'admin_manual_credit' }),
now: nowIso(),
})
publishWorkerWalletRealtimeChange(Number(workerId))
return { wallet: mapWallet(wallet) }
}
export async function listAdminWorkerWithdrawalAccounts(workerId: number | string) {
const worker = await getRequiredWorker(workerId)
const items = await listWorkerWithdrawalAccounts(worker.id)
return {
worker: mapWorkerUser(worker),
items: items.map(mapAdminWorkerWithdrawalAccount),
}
}
export async function saveAdminWorkerWithdrawalAccount(
workerId: number | string,
accountChannelValue: unknown,
payload: JsonObject = {},
) {
const worker = await getRequiredWorker(workerId)
const accountChannel = normalizeWithdrawChannel(accountChannelValue)
const currentAccounts = await listWorkerWithdrawalAccounts(worker.id)
const current = currentAccounts.find((item) => item.account_channel === accountChannel)
const accountName = String(payload.accountName || payload.realName || '').trim()
const accountNo =
accountChannel === 'alipay'
? String(payload.accountNo || payload.account || current?.account_no || '').trim()
: ''
const uploadedAlipayQrCode = normalizeProofFiles(
payload.alipayQrCodeImages || payload.alipayQrCodes,
)[0]
const alipayQrCodeImage =
accountChannel === 'alipay'
? uploadedAlipayQrCode || safeParseJson(current?.alipay_qr_code_json)
: {}
const uploadedWechatQrCode = normalizeProofFiles(
payload.wechatQrCodeImages || payload.wechatQrCodes,
)[0]
const wechatQrCodeImage =
accountChannel === 'wechat'
? uploadedWechatQrCode || safeParseJson(current?.wechat_qr_code_json)
: {}
if (!accountName) {
throw createHttpError('请填写收款人姓名', {
statusCode: 400,
errorCode: 'admin_worker_withdraw_account_name_required',
})
}
if (accountChannel === 'alipay' && !current && !accountNo) {
throw createHttpError('请填写支付宝账号', {
statusCode: 400,
errorCode: 'admin_worker_withdraw_alipay_account_required',
})
}
if (accountChannel === 'alipay' && !String(alipayQrCodeImage.url || '').trim()) {
throw createHttpError('请上传支付宝收款二维码', {
statusCode: 400,
errorCode: 'admin_worker_withdraw_alipay_qr_required',
})
}
if (accountChannel === 'wechat' && !String(wechatQrCodeImage.url || '').trim()) {
throw createHttpError('请上传微信收款二维码', {
statusCode: 400,
errorCode: 'admin_worker_withdraw_wechat_qr_required',
})
}
const saved = await upsertWorkerWithdrawalAccount({
workerId: worker.id,
accountChannel,
accountName,
accountNo,
alipayQrCodeJson: JSON.stringify(alipayQrCodeImage),
wechatQrCodeJson: JSON.stringify(wechatQrCodeImage),
now: nowIso(),
})
if (!saved) {
throw createHttpError('保存打手提现信息失败', {
statusCode: 500,
errorCode: 'admin_worker_withdraw_account_save_failed',
})
}
return { account: mapAdminWorkerWithdrawalAccount(saved) }
}
export async function getAdminWorkerFinanceConfig() {
return getWorkerFinanceConfig()
}
export async function saveAdminWorkerFinanceConfig(payload: JsonObject = {}) {
const config = await saveWorkerFinanceConfig(payload)
return config
}
export async function listAdminWorkerFinanceRequests(query: JsonObject = {}) {
const page = normalizePage(query.page)
const pageSize = normalizePageSize(query.pageSize)
const requestId = normalizeOptionalId(query.requestId ?? query.request_id) || 0
const status = normalizeFinanceRequestStatus(query.status)
const requestType = normalizeFinanceRequestType(query.requestType)
const keyword = String(query.keyword || '').trim()
const { items, total } = await listWorkerFinanceRequests({
page,
pageSize,
requestId,
status,
requestType,
keyword,
})
return {
items: items.map(mapFinanceRequest),
pagination: { page, pageSize, total },
}
}
export async function reviewAdminWorkerFinanceRequest(
requestId: number | string,
payload: JsonObject = {},
) {
const current = await getWorkerFinanceRequestById(requestId)
if (!current) {
throw createHttpError('资金申请不存在', {
statusCode: 404,
errorCode: 'worker_finance_request_not_found',
})
}
const status = normalizeAdminFinanceReviewStatus(payload.status)
const reviewedNote = String(payload.reviewedNote || payload.note || '').trim()
const reviewed = await reviewWorkerFinanceRequest({
requestId: Number(current.id),
status,
reviewedNote,
now: nowIso(),
})
if (reviewed.failureReason === 'request_not_pending') {
throw createHttpError('该申请已处理,请刷新后重试', {
statusCode: 409,
errorCode: 'worker_finance_request_reviewed',
})
}
if (reviewed.failureReason === 'withdraw_insufficient') {
throw createHttpError('打手余额不足,暂时无法通过该提现申请', {
statusCode: 409,
errorCode: 'worker_finance_request_withdraw_insufficient',
})
}
if (!reviewed.request) {
throw createHttpError('资金申请不存在', {
statusCode: 404,
errorCode: 'worker_finance_request_not_found',
})
}
await resolveAdminNotificationEntity(
'worker_finance_request',
Number(reviewed.request.id),
'finance_request_reviewed',
)
publishWorkerFinanceRealtimeChange({
requestId: Number(reviewed.request.id),
workerId: Number(reviewed.request.worker_id),
walletChanged: reviewed.request.request_type === 'withdraw' && status === 'approved',
})
return {
request: mapFinanceRequest(reviewed.request),
}
}
@@ -0,0 +1,287 @@
import { WORK_ORDER_STATUS } from '../../domain/work-order-status.js'
import {
createAfterSalesCase,
getAfterSalesCaseById,
listAfterSalesCaseEvents,
listAfterSalesCases,
listWorkerAfterSalesCases,
resolveAfterSalesCase,
submitWorkerAfterSalesResponse,
} from '../../repositories/worker-platform/after-sales-repo.js'
import { createWorkOrderEvent } from '../../repositories/worker-platform/work-order-event-repo.js'
import type { JsonObject } from '../../types/json.js'
import { createHttpError } from '../../utils/http.js'
import { randomId } from '../../utils/random.js'
import { nowIso } from '../../utils/time.js'
import { normalizePage, normalizePageSize, safeParseJson } from '../admin/admin-query-utils.js'
import {
publishWorkerWalletRealtimeChange,
publishWorkOrderRealtimeChange,
} from '../realtime/realtime-event-service.js'
import {
mapWorkerAfterSalesCase,
normalizeAmountFen,
normalizeOptionalId,
normalizeProofFiles,
normalizeUploadedFiles,
} from './mappers.js'
import { requireActiveWorkerSession, type WorkerSession } from './worker-service.js'
export async function listAdminAfterSalesCases(query: JsonObject = {}) {
const page = normalizePage(query.page)
const pageSize = normalizePageSize(query.pageSize)
const workOrderId = normalizeOptionalId(query.workOrderId ?? query.work_order_id)
const result = await listAfterSalesCases({
page,
pageSize,
status: String(query.status || '').trim(),
keyword: String(query.keyword || '').trim(),
...(workOrderId ? { workOrderId } : {}),
})
return {
items: result.items.map(mapWorkerAfterSalesCase),
pagination: { page, pageSize, total: result.total },
}
}
export async function getAdminAfterSalesCaseEvents(caseId: number | string) {
const afterSalesCase = await getAfterSalesCaseById(caseId)
if (!afterSalesCase) {
throw createHttpError('售后问题单不存在', {
statusCode: 404,
errorCode: 'after_sales_case_not_found',
})
}
return {
items: (await listAfterSalesCaseEvents(afterSalesCase.id)).map((event) => ({
eventId: Number(event.id),
actorType: event.actor_type,
actorId: event.actor_id,
eventType: event.event_type,
payload: safeParseJson(event.payload_json),
createdAt: event.created_at,
})),
}
}
export async function createAdminAfterSalesCase(
workOrderId: number | string,
payload: JsonObject = {},
actorName = '',
) {
const complaintNote = String(payload.complaintNote || payload.note || '').trim()
if (!complaintNote) {
throw createHttpError('请填写用户投诉说明', {
statusCode: 400,
errorCode: 'after_sales_case_complaint_required',
})
}
if (complaintNote.length > 2_000) {
throw createHttpError('用户投诉说明不能超过 2000 个字符', {
statusCode: 400,
errorCode: 'after_sales_case_complaint_too_long',
})
}
const problemType = String(payload.problemType || 'fake_evidence').trim() || 'fake_evidence'
if (problemType.length > 64) {
throw createHttpError('问题类型不正确', {
statusCode: 400,
errorCode: 'after_sales_case_problem_type_invalid',
})
}
const shareId = normalizeOptionalId(payload.shareId ?? payload.workOrderShareId)
const now = nowIso()
const result = await createAfterSalesCase({
caseNo: randomId('AS'),
workOrderId: Number(workOrderId),
workOrderShareId: shareId || null,
problemType,
complaintNote,
complaintFilesJson: JSON.stringify(normalizeProofFiles(payload.files || payload.proofFiles)),
createdBy: actorName,
now,
})
if (result.failureReason === 'work_order_not_accepted') {
throw createHttpError('只有已验收的工单可以发起售后问题单', {
statusCode: 409,
errorCode: 'after_sales_case_work_order_status_invalid',
})
}
if (result.failureReason === 'work_order_share_not_accepted') {
throw createHttpError('该拼单份额尚未验收,不能发起售后问题单', {
statusCode: 409,
errorCode: 'after_sales_case_share_status_invalid',
})
}
if (result.failureReason === 'worker_missing') {
throw createHttpError('请指定已验收的责任打手或拼单份额', {
statusCode: 409,
errorCode: 'after_sales_case_worker_required',
})
}
if (result.failureReason === 'active_case_exists' || !result.case) {
throw createHttpError('该责任打手已有处理中售后问题单', {
statusCode: 409,
errorCode: 'after_sales_case_active_exists',
})
}
await createWorkOrderEvent({
workOrderId: Number(result.case.work_order_id),
actorType: 'admin',
actorId: actorName,
eventType: 'after_sales_case_created',
fromStatus: WORK_ORDER_STATUS.ACCEPTED,
toStatus: WORK_ORDER_STATUS.ACCEPTED,
payloadJson: JSON.stringify({
caseId: Number(result.case.id),
caseNo: result.case.case_no,
workerId: Number(result.case.worker_id),
shareId: result.case.work_order_share_id ? Number(result.case.work_order_share_id) : null,
problemType,
}),
now,
})
publishWorkOrderRealtimeChange({
workOrderId: Number(result.case.work_order_id),
workerIds: [Number(result.case.worker_id)],
})
return { case: mapWorkerAfterSalesCase(result.case) }
}
export async function resolveAdminAfterSalesCase(
caseId: number | string,
payload: JsonObject = {},
actorName = '',
) {
const action = String(payload.action || '').trim()
if (!['not_upheld', 'remedied', 'upheld'].includes(action)) {
throw createHttpError('售后处置方式不正确', {
statusCode: 400,
errorCode: 'after_sales_case_resolution_action_invalid',
})
}
const resolutionNote = String(payload.resolutionNote || payload.note || '').trim()
if (!resolutionNote) {
throw createHttpError('请填写售后处置备注', {
statusCode: 400,
errorCode: 'after_sales_case_resolution_note_required',
})
}
if (resolutionNote.length > 2_000) {
throw createHttpError('售后处置备注不能超过 2000 个字符', {
statusCode: 400,
errorCode: 'after_sales_case_resolution_note_too_long',
})
}
const recoveryAmount = normalizeAmountFen(payload.recoveryAmount ?? payload.amount, 0)
if (action === 'upheld' && recoveryAmount <= 0) {
throw createHttpError('问题成立时请填写大于 0 的追缴金额', {
statusCode: 400,
errorCode: 'after_sales_case_recovery_amount_required',
})
}
const result = await resolveAfterSalesCase({
caseId: Number(caseId),
action: action as 'not_upheld' | 'remedied' | 'upheld',
resolutionNote,
recoveryAmount: action === 'upheld' ? recoveryAmount : 0,
resolvedBy: actorName,
now: nowIso(),
})
if (result.failureReason === 'not_found') {
throw createHttpError('售后问题单不存在', {
statusCode: 404,
errorCode: 'after_sales_case_not_found',
})
}
if (result.failureReason === 'not_resolvable' || !result.case) {
throw createHttpError('该售后问题单已处置,请刷新后重试', {
statusCode: 409,
errorCode: 'after_sales_case_not_resolvable',
})
}
await createWorkOrderEvent({
workOrderId: Number(result.case.work_order_id),
actorType: 'admin',
actorId: actorName,
eventType: 'after_sales_case_resolved',
fromStatus: WORK_ORDER_STATUS.ACCEPTED,
toStatus: WORK_ORDER_STATUS.ACCEPTED,
payloadJson: JSON.stringify({
caseId: Number(result.case.id),
action,
recoveryAmount: Number(result.case.recovery_amount || 0),
debtAmount: Number(result.case.debt_amount || 0),
}),
now: nowIso(),
})
const workerId = Number(result.case.worker_id)
publishWorkOrderRealtimeChange({
workOrderId: Number(result.case.work_order_id),
workerIds: workerId > 0 ? [workerId] : [],
})
if (workerId > 0) publishWorkerWalletRealtimeChange(workerId)
return { case: mapWorkerAfterSalesCase(result.case) }
}
export async function getWorkerAfterSalesCaseSummary(session: WorkerSession) {
requireActiveWorkerSession(session)
return { items: (await listWorkerAfterSalesCases(session.workerId)).map(mapWorkerAfterSalesCase) }
}
export async function submitWorkerAfterSalesCaseResponse(
caseId: number | string,
payload: JsonObject = {},
session: WorkerSession,
) {
requireActiveWorkerSession(session)
const response = String(payload.response || payload.note || '').trim()
const files = normalizeUploadedFiles(payload.files || payload.proofFiles)
if (!response && files.length === 0) {
throw createHttpError('请填写说明或上传补救材料', {
statusCode: 400,
errorCode: 'after_sales_case_response_required',
})
}
if (response.length > 2_000) {
throw createHttpError('售后说明不能超过 2000 个字符', {
statusCode: 400,
errorCode: 'after_sales_case_response_too_long',
})
}
const now = nowIso()
const result = await submitWorkerAfterSalesResponse({
caseId: Number(caseId),
workerId: Number(session.workerId),
response,
filesJson: JSON.stringify(files),
now,
})
if (result.failureReason === 'not_found') {
throw createHttpError('售后问题单不存在或不属于当前账号', {
statusCode: 404,
errorCode: 'after_sales_case_not_found',
})
}
if (result.failureReason === 'not_pending' || !result.case) {
throw createHttpError('该售后问题单已提交说明或已处置', {
statusCode: 409,
errorCode: 'after_sales_case_response_status_invalid',
})
}
await createWorkOrderEvent({
workOrderId: Number(result.case.work_order_id),
actorType: 'worker',
actorId: String(session.workerId),
eventType: 'after_sales_case_responded',
fromStatus: WORK_ORDER_STATUS.ACCEPTED,
toStatus: WORK_ORDER_STATUS.ACCEPTED,
payloadJson: JSON.stringify({ caseId: Number(result.case.id), fileCount: files.length }),
now,
})
publishWorkOrderRealtimeChange({
workOrderId: Number(result.case.work_order_id),
workerIds: [Number(session.workerId)],
})
return { case: mapWorkerAfterSalesCase(result.case) }
}
@@ -1,6 +1,11 @@
export * from './mappers.js'
export * from './worker-service.js'
export * from './after-sales-service.js'
export * from './admin-service.js'
export * from './admin-worker-finance-service.js'
export * from './admin-worker-catalog-service.js'
export * from './admin-kuaishou-match-service.js'
export * from './admin-worker-config-service.js'
export * from './worker-product-match-config-service.js'
export * from './worker-announcement-config-service.js'
export * from './worker-hall-config-service.js'
@@ -33,7 +33,7 @@ import {
updateWorkOrder,
} from '../../repositories/worker-platform/work-order-repo.js'
import { syncWorkerOrdersForSourceOrder } from './admin-service.js'
import { createWorkOrderEvent } from '../../repositories/worker-platform/work-order-repo.js'
import { createWorkOrderEvent } from '../../repositories/worker-platform/work-order-event-repo.js'
import { OPEN_91_PLATFORM, OPEN_91_PROVIDER } from '../open-91/config.js'
import {
findKuaishouIndustryShopConfig,
@@ -8,7 +8,6 @@ import {
countWorkerAcceptedOrders,
createWorkerCancelRequest,
createWorkerOrderFeedback,
listWorkerAfterSalesCases,
countWorkerActiveOrders,
countWorkerTimeoutEvents,
countWorkerWithdrawRequestsOnDay,
@@ -65,7 +64,6 @@ import {
releaseDepositUnfreeze,
settleOverdueWorkOrder,
submitWorkOrderShareAcceptance,
submitWorkerAfterSalesResponse,
updateWorkOrder,
updateWorkOrderAcceptanceRecord,
updateWorkOrderDraftAcceptance,
@@ -124,7 +122,6 @@ import {
mapWorkOrderEvents,
mapWorkOrderShare,
mapWorkerCancelRequest,
mapWorkerAfterSalesCase,
mapWorkerOrderFeedback,
mapWorkerUser,
normalizeAmountFen,
@@ -1858,68 +1855,6 @@ export async function getWorkerOrderFeedbackSummary(session: WorkerSession) {
return { items: (await listWorkerOrderFeedbacks(session.workerId)).map(mapWorkerOrderFeedback) }
}
export async function getWorkerAfterSalesCaseSummary(session: WorkerSession) {
requireActiveWorkerSession(session)
return { items: (await listWorkerAfterSalesCases(session.workerId)).map(mapWorkerAfterSalesCase) }
}
export async function submitWorkerAfterSalesCaseResponse(
caseId: number | string,
payload: JsonObject = {},
session: WorkerSession,
) {
requireActiveWorkerSession(session)
const response = String(payload.response || payload.note || '').trim()
const files = normalizeUploadedFiles(payload.files || payload.proofFiles)
if (!response && files.length === 0) {
throw createHttpError('请填写说明或上传补救材料', {
statusCode: 400,
errorCode: 'after_sales_case_response_required',
})
}
if (response.length > 2_000) {
throw createHttpError('售后说明不能超过 2000 个字符', {
statusCode: 400,
errorCode: 'after_sales_case_response_too_long',
})
}
const now = nowIso()
const result = await submitWorkerAfterSalesResponse({
caseId: Number(caseId),
workerId: Number(session.workerId),
response,
filesJson: JSON.stringify(files),
now,
})
if (result.failureReason === 'not_found') {
throw createHttpError('售后问题单不存在或不属于当前账号', {
statusCode: 404,
errorCode: 'after_sales_case_not_found',
})
}
if (result.failureReason === 'not_pending' || !result.case) {
throw createHttpError('该售后问题单已提交说明或已处置', {
statusCode: 409,
errorCode: 'after_sales_case_response_status_invalid',
})
}
await createWorkOrderEvent({
workOrderId: Number(result.case.work_order_id),
actorType: 'worker',
actorId: String(session.workerId),
eventType: 'after_sales_case_responded',
fromStatus: WORK_ORDER_STATUS.ACCEPTED,
toStatus: WORK_ORDER_STATUS.ACCEPTED,
payloadJson: JSON.stringify({ caseId: Number(result.case.id), fileCount: files.length }),
now,
})
publishWorkOrderRealtimeChange({
workOrderId: Number(result.case.work_order_id),
workerIds: [Number(session.workerId)],
})
return { case: mapWorkerAfterSalesCase(result.case) }
}
export async function submitWorkerOrderFeedback(
workOrderId: number | string,
payload: JsonObject = {},