diff --git a/apps/backend/src/services/worker-platform/index.ts b/apps/backend/src/services/worker-platform/index.ts index 49e59b07..77924462 100644 --- a/apps/backend/src/services/worker-platform/index.ts +++ b/apps/backend/src/services/worker-platform/index.ts @@ -10,6 +10,8 @@ export * from './worker-registration-service.js' export * from './worker-login-service.js' export * from './worker-profile-service.js' export * from './worker-hall-service.js' +export * from './worker-order-timeout-service.js' +export * from './worker-order-realtime-service.js' export * from './worker-platform-defaults.js' export * from './work-order-timeout-policy.js' export * from './after-sales-service.js' diff --git a/apps/backend/src/services/worker-platform/worker-order-realtime-service.ts b/apps/backend/src/services/worker-platform/worker-order-realtime-service.ts new file mode 100644 index 00000000..6d11d5e5 --- /dev/null +++ b/apps/backend/src/services/worker-platform/worker-order-realtime-service.ts @@ -0,0 +1,32 @@ +import { listWorkOrderShares } from '../../repositories/worker-platform/index.js' +import { + publishWorkerWalletRealtimeChange, + publishWorkOrderRealtimeChange, +} from '../realtime/realtime-event-service.js' + +/** 拼单工单同步所有参与打手,避免仅主打手页面更新。 */ +export async function publishWorkOrderParticipantChange( + workOrderId: number, + options: { + workerIds?: number[] + hallChanged?: boolean + walletChangedWorkerIds?: number[] + } = {}, +) { + const shares = await listWorkOrderShares(workOrderId) + const workerIds = [ + ...(options.workerIds || []), + ...shares.map((share) => Number(share.worker_id)), + ].filter( + (workerId, index, values) => + Number.isInteger(workerId) && workerId > 0 && values.indexOf(workerId) === index, + ) + publishWorkOrderRealtimeChange({ + workOrderId, + workerIds, + ...(options.hallChanged ? { hallChanged: true } : {}), + }) + for (const workerId of options.walletChangedWorkerIds || []) { + publishWorkerWalletRealtimeChange(workerId) + } +} diff --git a/apps/backend/src/services/worker-platform/worker-order-timeout-service.ts b/apps/backend/src/services/worker-platform/worker-order-timeout-service.ts new file mode 100644 index 00000000..53f40f82 --- /dev/null +++ b/apps/backend/src/services/worker-platform/worker-order-timeout-service.ts @@ -0,0 +1,47 @@ +import { WORK_ORDER_STATUS } from '../../domain/work-order-status.js' +import { + listOverdueWorkOrders, + settleOverdueWorkOrder, +} from '../../repositories/worker-platform/index.js' +import { nowIso } from '../../utils/time.js' +import { + publishWorkerWalletRealtimeChange, + publishWorkOrderRealtimeChange, +} from '../realtime/realtime-event-service.js' +import { normalizeWorkOrderTimeoutPolicy } from './work-order-timeout-policy.js' + +export async function settleOverdueWorkOrders(options: { limit?: number; workerId?: number } = {}) { + const overdue = await listOverdueWorkOrders({ + limit: Math.max(1, Number(options.limit || 50)), + workerId: Number(options.workerId || 0), + }) + let processedCount = 0 + let skippedCount = 0 + for (const workOrder of overdue) { + const result = await settleOverdueWorkOrder({ + workOrderId: workOrder.id, + policy: normalizeWorkOrderTimeoutPolicy(workOrder.timeout_policy), + now: nowIso(), + }) + if (result.order) { + processedCount += 1 + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + workerIds: Number(workOrder.assigned_worker_id) + ? [Number(workOrder.assigned_worker_id)] + : [], + hallChanged: result.order.status === WORK_ORDER_STATUS.OPEN, + }) + if (Number(workOrder.assigned_worker_id)) { + publishWorkerWalletRealtimeChange(Number(workOrder.assigned_worker_id)) + } + } else { + skippedCount += 1 + } + } + return { + checkedCount: overdue.length, + processedCount, + skippedCount, + } +} diff --git a/apps/backend/src/services/worker-platform/worker-service.ts b/apps/backend/src/services/worker-platform/worker-service.ts index 2be44834..4d5892a0 100644 --- a/apps/backend/src/services/worker-platform/worker-service.ts +++ b/apps/backend/src/services/worker-platform/worker-service.ts @@ -11,6 +11,7 @@ import { createWorkOrderEvent, getLatestWorkerSmsCode, recordWorkerSmsCodeFailure, + settleOverdueWorkOrder, getWorkerUserById, getWorkerUserByPhone, getWorkerUserByUsername, @@ -33,10 +34,7 @@ import { listWorkerWorkOrderViews, listWorkOrders, listWorkOrderEventsByOrderId, - listWorkOrderShares, getWorkerCancelRequestRisk, - listOverdueWorkOrders, - settleOverdueWorkOrder, submitWorkOrderShareAcceptance, updateWorkOrder, updateWorkOrderAcceptanceRecord, @@ -123,6 +121,8 @@ import { verifyWorkerSessionToken, } from './worker-session-auth-service.js' import { normalizeWorkOrderTimeoutPolicy } from './work-order-timeout-policy.js' +import { publishWorkOrderParticipantChange } from './worker-order-realtime-service.js' +import { settleOverdueWorkOrders } from './worker-order-timeout-service.js' const VIP_AUTO_ACCEPT_EVIDENCE_WINDOW_HOURS = 24 @@ -278,41 +278,7 @@ function resolveWorkOrderDeadlineAt(workOrder: WorkOrderRow, now: string): strin export { settleDueDepositUnfreezes } from './worker-profile-service.js' -export async function settleOverdueWorkOrders(options: { limit?: number; workerId?: number } = {}) { - const overdue = await listOverdueWorkOrders({ - limit: Math.max(1, Number(options.limit || 50)), - workerId: Number(options.workerId || 0), - }) - let processedCount = 0 - let skippedCount = 0 - for (const workOrder of overdue) { - const result = await settleOverdueWorkOrder({ - workOrderId: workOrder.id, - policy: normalizeWorkOrderTimeoutPolicy(workOrder.timeout_policy), - now: nowIso(), - }) - if (result.order) { - processedCount += 1 - publishWorkOrderRealtimeChange({ - workOrderId: Number(workOrder.id), - workerIds: Number(workOrder.assigned_worker_id) - ? [Number(workOrder.assigned_worker_id)] - : [], - hallChanged: result.order.status === WORK_ORDER_STATUS.OPEN, - }) - if (Number(workOrder.assigned_worker_id)) { - publishWorkerWalletRealtimeChange(Number(workOrder.assigned_worker_id)) - } - } else { - skippedCount += 1 - } - } - return { - checkedCount: overdue.length, - processedCount, - skippedCount, - } -} +export { settleOverdueWorkOrders } from './worker-order-timeout-service.js' export async function listWorkerMyOrders(query: JsonObject = {}, session: WorkerSession) { requireActiveWorkerSession(session) @@ -1497,33 +1463,6 @@ export async function collectSubmitWorkOrder(payload: JsonObject = {}) { } } -/** 拼单工单同步所有参与打手,避免仅主打手页面更新。 */ -async function publishWorkOrderParticipantChange( - workOrderId: number, - options: { - workerIds?: number[] - hallChanged?: boolean - walletChangedWorkerIds?: number[] - } = {}, -) { - const shares = await listWorkOrderShares(workOrderId) - const workerIds = [ - ...(options.workerIds || []), - ...shares.map((share) => Number(share.worker_id)), - ].filter( - (workerId, index, values) => - Number.isInteger(workerId) && workerId > 0 && values.indexOf(workerId) === index, - ) - publishWorkOrderRealtimeChange({ - workOrderId, - workerIds, - ...(options.hallChanged ? { hallChanged: true } : {}), - }) - for (const workerId of options.walletChangedWorkerIds || []) { - publishWorkerWalletRealtimeChange(workerId) - } -} - export function resolveCollectSubmitTargetWorkOrder( orderNo: string, workOrders: WorkOrderRow[],