接单工单支持任务时限,超时自动判定失败并处置

This commit is contained in:
yml2213
2026-08-01 19:02:57 +08:00
parent 66193c5d59
commit 23ad60a1ef
18 changed files with 681 additions and 61 deletions
@@ -10,6 +10,7 @@ import {
const SCHEDULED_JOBS_FILE_PATH = path.join(PROJECT_ROOT, 'data', 'scheduled-jobs.json')
const CLOUDTENTACLES_HEALTH_JOB_ID = 'cloudtentacles-health'
const WORK_ORDER_TIMEOUT_JOB_ID = 'work-order-timeout'
export function getScheduledJobsFilePath() {
return SCHEDULED_JOBS_FILE_PATH
@@ -39,17 +40,32 @@ export function getCloudtentaclesHealthJob(config: JsonObject = getScheduledJobs
|| createDefaultCloudtentaclesHealthJob()
}
export function getWorkOrderTimeoutJob(config: JsonObject = getScheduledJobsConfig()) {
return (Array.isArray(config.jobs) ? config.jobs : [])
.find((item) => String(item.id || '').trim() === WORK_ORDER_TIMEOUT_JOB_ID)
|| createDefaultWorkOrderTimeoutJob()
}
export function getScheduledJobById(jobId: unknown, config: JsonObject = getScheduledJobsConfig()) {
return (Array.isArray(config.jobs) ? config.jobs : [])
.find((item) => String(item.id || '').trim() === String(jobId || '').trim()) || null
}
export function normalizeScheduledJobsConfig(rawValue: unknown) {
const source = isPlainObject(rawValue) ? rawValue : {}
const rawJobs = Array.isArray(source.jobs) ? source.jobs : []
const jobs = rawJobs
.map((item) => normalizeScheduledJob(item))
.filter((item): item is ReturnType<typeof normalizeCloudtentaclesHealthJob> => Boolean(item))
.filter((item): item is NonNullable<ReturnType<typeof normalizeScheduledJob>> => Boolean(item))
const hasCloudtentaclesHealth = jobs.some((item) => item.id === CLOUDTENTACLES_HEALTH_JOB_ID)
const hasWorkOrderTimeout = jobs.some((item) => item.id === WORK_ORDER_TIMEOUT_JOB_ID)
if (!hasCloudtentaclesHealth) {
jobs.push(createDefaultCloudtentaclesHealthJob())
}
if (!hasWorkOrderTimeout) {
jobs.push(createDefaultWorkOrderTimeoutJob())
}
return {
enabled: typeof source.enabled === 'boolean' ? source.enabled : true,
@@ -63,11 +79,28 @@ function normalizeScheduledJob(rawValue: unknown) {
}
const type = String(rawValue.type || '').trim()
if (type !== 'cloudtentacles_health') {
return null
if (type === 'cloudtentacles_health') {
return normalizeCloudtentaclesHealthJob(rawValue)
}
if (type === 'work_order_timeout') {
return normalizeWorkOrderTimeoutJob(rawValue)
}
return normalizeCloudtentaclesHealthJob(rawValue)
return null
}
function normalizeWorkOrderTimeoutJob(rawValue: JsonObject) {
const config = isPlainObject(rawValue.config) ? rawValue.config : {}
return {
id: WORK_ORDER_TIMEOUT_JOB_ID,
type: 'work_order_timeout',
enabled: rawValue.enabled !== false,
intervalSeconds: normalizeRangeInteger(rawValue.intervalSeconds, 60, 30, 3600),
config: {
scanLimit: normalizeRangeInteger(config.scanLimit, 50, 1, 1000),
},
}
}
function normalizeCloudtentaclesHealthJob(rawValue: JsonObject) {
@@ -137,10 +170,23 @@ function createDefaultScheduledJobsConfig() {
enabled: true,
jobs: [
createDefaultCloudtentaclesHealthJob(),
createDefaultWorkOrderTimeoutJob(),
],
}
}
function createDefaultWorkOrderTimeoutJob() {
return {
id: WORK_ORDER_TIMEOUT_JOB_ID,
type: 'work_order_timeout',
enabled: true,
intervalSeconds: 60,
config: {
scanLimit: 50,
},
}
}
function createDefaultCloudtentaclesHealthJob() {
return {
id: CLOUDTENTACLES_HEALTH_JOB_ID,
@@ -2,10 +2,11 @@ import { logError, logInfo, logWarn } from '../../utils/logger.js'
import { createHttpError } from '../../utils/http.js'
import type { JsonObject } from '../../types/json.js'
import {
getCloudtentaclesHealthJob,
getScheduledJobById,
getScheduledJobsConfig,
} from './config-service.js'
import { runCloudtentaclesHealthJob } from './cloudtentacles-health-job.js'
import { runWorkOrderTimeoutJob } from './work-order-timeout-job.js'
const timers = new Map<string, NodeJS.Timeout>()
const jobStates = new Map<string, JsonObject>()
@@ -67,14 +68,20 @@ export async function runScheduledJobNow(jobId: unknown) {
const result = await runJob(job, { manual: true })
if (getScheduledJobsConfig().enabled !== false) {
const latestJob = getCloudtentaclesHealthJob(getScheduledJobsConfig())
if (latestJob.enabled === true) {
scheduleJob(latestJob, Math.max(60, Number(latestJob.intervalSeconds || 300)) * 1000)
}
rescheduleJob(job)
}
return result
}
function rescheduleJob(job: JsonObject) {
const jobId = String(job.id || '').trim()
if (!jobId) return
const latestJob = getScheduledJobById(jobId)
if (!latestJob || latestJob.enabled !== true) return
const intervalMs = Math.max(60, Number(latestJob.intervalSeconds || 60)) * 1000
scheduleJob(latestJob, intervalMs)
}
function scheduleJob(job: JsonObject, delayMs: number) {
const jobId = String(job.id || '').trim()
if (!jobId) {
@@ -90,10 +97,7 @@ function scheduleJob(job: JsonObject, delayMs: number) {
const timer = setTimeout(() => {
timers.delete(jobId)
void runJob(job).finally(() => {
const latestJob = getCloudtentaclesHealthJob(getScheduledJobsConfig())
if (latestJob.enabled === true) {
scheduleJob(latestJob, Math.max(60, Number(latestJob.intervalSeconds || 300)) * 1000)
}
rescheduleJob(job)
})
}, delayMs)
@@ -126,25 +130,27 @@ async function runJob(job: JsonObject, { manual = false }: { manual?: boolean }
try {
const result = await dispatchJob(job)
const summary =
result && typeof result === 'object' ? (result as JsonObject) : {}
updateJobState(jobId, {
running: false,
lastFinishedAt: new Date().toISOString(),
lastStatus: result?.status || 'ok',
lastMessage: result?.message || '执行完成',
lastAsset: typeof result?.asset === 'number' ? result.asset : null,
lastThreshold: typeof result?.threshold === 'number' ? result.threshold : null,
lastAccounts: Array.isArray(result?.accounts) ? result.accounts : [],
lastAccountCount: Number(result?.accountCount || 0),
lastCheckedCount: Number(result?.checkedCount || 0),
lastOkCount: Number(result?.okCount || 0),
lastLowAssetCount: Number(result?.lowAssetCount || 0),
lastFailedCount: Number(result?.failedCount || 0),
lastStatus: String(summary.status || 'ok'),
lastMessage: String(summary.message || '执行完成'),
lastAsset: typeof summary.asset === 'number' ? summary.asset : null,
lastThreshold: typeof summary.threshold === 'number' ? summary.threshold : null,
lastAccounts: Array.isArray(summary.accounts) ? summary.accounts : [],
lastAccountCount: Number(summary.accountCount || 0),
lastCheckedCount: Number(summary.checkedCount || 0),
lastOkCount: Number(summary.okCount || 0),
lastLowAssetCount: Number(summary.lowAssetCount || 0),
lastFailedCount: Number(summary.failedCount || 0),
lastManual: manual,
})
logInfo('[scheduler]', '定时任务执行完成', {
jobId,
type: job.type,
status: result?.status || 'ok',
status: String(summary.status || 'ok'),
})
return result
} catch (error) {
@@ -173,6 +179,9 @@ function dispatchJob(job: JsonObject) {
if (job.type === 'cloudtentacles_health') {
return runCloudtentaclesHealthJob(job)
}
if (job.type === 'work_order_timeout') {
return runWorkOrderTimeoutJob(job)
}
throw createHttpError(`不支持的定时任务类型:${job.type}`, {
statusCode: 400,
@@ -0,0 +1,24 @@
import type { JsonObject } from '../../types/json.js'
import { logInfo } from '../../utils/logger.js'
import { settleOverdueWorkOrders } from '../worker-platform/worker-service.js'
export async function runWorkOrderTimeoutJob(job: JsonObject) {
const config =
job.config && typeof job.config === 'object' && !Array.isArray(job.config)
? (job.config as JsonObject)
: {}
const scanLimit = Math.max(1, Number(config.scanLimit || 50))
const result = await settleOverdueWorkOrders({ limit: scanLimit })
logInfo('[work-order-timeout]', '超时工单扫描完成', result)
return {
ok: result.processedCount === result.checkedCount,
status: 'ok',
message: `扫描 ${result.checkedCount} 个,处置 ${result.processedCount} 个,跳过 ${result.skippedCount}`,
checkedCount: result.checkedCount,
failedCount: result.skippedCount,
asset: result.processedCount,
}
}
@@ -86,7 +86,7 @@ import {
} from './worker-finance-config-service.js'
import { DEFAULT_CATEGORY_KEY, DEFAULT_DEPOSIT_THRESHOLD_AMOUNT, DEFAULT_LEVEL_KEY, DEFAULT_LEVEL_NAME, mapFinanceRequest, mapWallet, mapWorkCategory, mapWorkOrderAdmin, mapWorkOrderShare, mapWorkProductRule, mapWorkerLevel, mapWorkerUser, normalizeAdminFinanceReviewStatus, normalizeAmountFen, normalizeBoolean, normalizeEnabledStatus, normalizeFinanceRequestStatus, normalizeFinanceRequestType, normalizeInteger, normalizeMatchType, normalizeOptionalId, normalizePositiveInteger, normalizeProblemResolutionAction, normalizeRequirementFields, normalizeRequirementFieldsFromPayload, normalizeReviewStatus, normalizeSessionVersion, normalizeSlugKey, normalizeSubmittedFields, normalizeUploadedFiles, resolveMatchingProductRule, resolveRequirementFields } from './mappers.js'
import { ensureWorkerPlatformDefaults, getRequiredWorkOrder, getRequiredWorker } from './worker-service.js'
import { ensureWorkerPlatformDefaults, getRequiredWorkOrder, getRequiredWorker, normalizeWorkOrderTimeoutPolicy } from './worker-service.js'
export async function listAdminWorkerLevels() {
await ensureWorkerPlatformDefaults()
@@ -281,6 +281,8 @@ export async function saveAdminWorkProductRule(payload: JsonObject = {}) {
sharingEnabled,
sharingTotalQuantity,
sharingUnitReward,
timeoutMinutes: Math.max(0, normalizeInteger(payload.timeoutMinutes, 0)),
timeoutPolicy: normalizeWorkOrderTimeoutPolicy(payload.timeoutPolicy),
requirementJson: JSON.stringify({ fields }),
sortOrder: normalizeInteger(payload.sortOrder, 100),
now: nowIso(),
@@ -646,6 +648,8 @@ export async function createAdminMockWorkOrder(payload: JsonObject = {}) {
rewardAmount,
requiredDepositAmount,
depositThresholdAmount,
timeoutMinutes: Math.max(0, normalizeInteger(payload.timeoutMinutes, 0)),
timeoutPolicy: normalizeWorkOrderTimeoutPolicy(payload.timeoutPolicy),
materialJson: JSON.stringify(material),
requirementJson: JSON.stringify({ fields }),
now,
@@ -721,6 +725,13 @@ export async function updateAdminWorkOrder(
rewardAmount,
requiredDepositAmount,
depositThresholdAmount,
timeoutMinutes:
payload.timeoutMinutes === undefined
? workOrder.timeout_minutes
: Math.max(0, normalizeInteger(payload.timeoutMinutes, workOrder.timeout_minutes)),
timeoutPolicy: normalizeWorkOrderTimeoutPolicy(
payload.timeoutPolicy ?? workOrder.timeout_policy,
),
requirementJson: JSON.stringify({ fields }),
updatedAt: now,
})
@@ -735,6 +746,8 @@ export async function updateAdminWorkOrder(
productName: updated?.product_name,
rewardAmount,
requiredDepositAmount,
timeoutMinutes: updated?.timeout_minutes,
timeoutPolicy: updated?.timeout_policy,
fields,
}),
now,
@@ -995,6 +1008,8 @@ export async function syncWorkerOrdersForSourceOrder(
sharingEnabled: rule.sharing_enabled === true,
sharingTotalQuantity: Number(rule.sharing_total_quantity || 1),
sharingUnitReward: Number(rule.sharing_unit_reward || 0),
timeoutMinutes: Number(rule.timeout_minutes || 0),
timeoutPolicy: String(rule.timeout_policy || 'reopen').trim(),
materialJson: JSON.stringify({
source: {
orderId: Number(order.id),
@@ -247,6 +247,8 @@ export function mapWorkProductRule(rule: WorkProductRuleRow | null | undefined)
totalQuantity: Number(rule.sharing_total_quantity || 1),
unitReward: Number(rule.sharing_unit_reward || 0),
},
timeoutMinutes: Number(rule.timeout_minutes || 0),
timeoutPolicy: String(rule.timeout_policy || 'reopen').trim(),
requirement: {
fields: normalizeRequirementFields(requirement.fields),
},
@@ -355,6 +357,9 @@ export function mapWorkOrderAdmin(workOrder: WorkOrderRow) {
totalQuantity: Number(workOrder.sharing_total_quantity || 1),
unitReward: Number(workOrder.sharing_unit_reward || 0),
},
timeoutMinutes: Number(workOrder.timeout_minutes || 0),
timeoutPolicy: String(workOrder.timeout_policy || 'reopen').trim(),
deadlineAt: workOrder.deadline_at,
depositThresholdAmount: Number(
workOrder.deposit_threshold_amount || DEFAULT_DEPOSIT_THRESHOLD_AMOUNT,
),
@@ -47,8 +47,10 @@ import {
listWorkerUsers,
listWorkOrderShares,
listWorkOrderSharesByOrderIds,
listOverdueWorkOrders,
resolveProblemWorkOrder,
reviewWorkerFinanceRequest,
settleOverdueWorkOrder,
submitWorkOrderShareAcceptance,
updateWorkOrder,
updateWorkerPassword,
@@ -642,6 +644,7 @@ export async function listWorkerHallOrders(query: JsonObject = {}, session: Work
const pageSize = normalizePageSize(query.pageSize)
const categoryId = normalizeOptionalId(query.categoryId)
const worker = await getRequiredWorker(session.workerId)
await settleOverdueWorkOrders({ workerId: worker.id, limit: 50 })
const permissions = resolveWorkerPermissions(worker)
const visibleDelaySeconds = Number(permissions.visibleDelaySeconds || 0)
const visibleAfterIso =
@@ -706,6 +709,7 @@ export async function grabWorkerHallOrder(workOrderId: number | string, session:
workerId: Number(worker.id),
depositAmount: freezeAmount,
maxActiveOrders: permissions.maxActiveOrders,
deadlineAt: resolveWorkOrderDeadlineAt(workOrder, nowIso()),
now: nowIso(),
})
if (!grabbed.order) {
@@ -715,11 +719,61 @@ export async function grabWorkerHallOrder(workOrderId: number | string, session:
return { order: mapWorkOrderForWorker(grabbed.order, permissions) }
}
function resolveWorkOrderDeadlineAt(
workOrder: WorkOrderRow,
now: string,
): string | null {
const timeoutMinutes = Number(workOrder.timeout_minutes || 0)
if (timeoutMinutes <= 0) return null
return new Date(new Date(now).getTime() + timeoutMinutes * 60 * 1000).toISOString()
}
export const WORK_ORDER_TIMEOUT_POLICIES = ['reopen', 'cancel_release', 'cancel_deduct'] as const
export type WorkOrderTimeoutPolicy = (typeof WORK_ORDER_TIMEOUT_POLICIES)[number]
export function isWorkOrderTimeoutPolicy(value: unknown): value is WorkOrderTimeoutPolicy {
return WORK_ORDER_TIMEOUT_POLICIES.includes(value as WorkOrderTimeoutPolicy)
}
export function normalizeWorkOrderTimeoutPolicy(value: unknown): WorkOrderTimeoutPolicy {
return isWorkOrderTimeoutPolicy(value) ? value : 'reopen'
}
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
} else {
skippedCount += 1
}
}
return {
checkedCount: overdue.length,
processedCount,
skippedCount,
}
}
export async function listWorkerMyOrders(query: JsonObject = {}, session: WorkerSession) {
requireActiveWorkerSession(session)
const worker = await getRequiredWorker(session.workerId)
await settleOverdueWorkOrders({ workerId: worker.id, limit: 50 })
const page = normalizePage(query.page)
const pageSize = normalizePageSize(query.pageSize)
const worker = await getRequiredWorker(session.workerId)
const { items, total } = await listWorkOrders({
page,
pageSize,
@@ -835,6 +889,22 @@ export async function submitWorkerOrderAcceptance(
if (sharingShare) {
return submitWorkerSharingAcceptance(workOrder, sharingShare, payload, session)
}
const now = nowIso()
if (
workOrder.status === WORK_ORDER_STATUS.IN_PROGRESS &&
workOrder.deadline_at &&
new Date(workOrder.deadline_at).getTime() <= Date.now()
) {
await settleOverdueWorkOrder({
workOrderId: workOrder.id,
policy: normalizeWorkOrderTimeoutPolicy(workOrder.timeout_policy),
now,
})
throw createHttpError('任务已超时,系统已按超时策略处置', {
statusCode: 409,
errorCode: 'work_order_timeout_expired',
})
}
if (Number(workOrder.assigned_worker_id || 0) !== workerId) {
throw createHttpError('只能提交自己的订单', {
statusCode: 403,
@@ -850,7 +920,6 @@ export async function submitWorkerOrderAcceptance(
})
}
const now = nowIso()
const files = normalizeUploadedFiles(payload.files)
const imageUrls = [
...files.map((file) => file.url || file.mediumUrl || file.thumbnailUrl).filter(Boolean),