e60bd8567a
新增纯 TS 任务轮询协调器 task-polling-coordinator:重复 add 去重、 并发 runOnce 合并为单次请求并共享结果、markTerminal 幂等终态清理 (只回调一次)、maxTasks 上限拒绝超出条目。useTaskProgressLoop 新增 可选 coordinator 选项,add/remove/终态路径委托协调器统一处理(默认 不启用)。测试用 node:test 覆盖 8 个场景,含 3 路并发合并断言。
142 lines
3.7 KiB
TypeScript
142 lines
3.7 KiB
TypeScript
/**
|
|
* 任务轮询协调器(Task 83)。
|
|
*
|
|
* 统一不同页面的任务轮询语义:
|
|
* - 去重:同一 taskId 重复 add 只保留一份,不产生重复记录与重复请求;
|
|
* - in-flight 合并:并发 runOnce 合并为单次请求,全部调用共享同一次结果;
|
|
* - 终态清理:markTerminal 从集合移除任务并回调 onTerminal,重复调用幂等
|
|
* (只回调一次),终态后任务不再被轮询;
|
|
* - 有界:maxTasks 上限内的 add 被拒绝并计数,防止任务集无界增长。
|
|
*
|
|
* 纯 TS 无副作用模块;add 对非法 taskId fail-fast 抛错,loader 抛错时
|
|
* in-flight 释放且任务集合保持不变,可恢复后继续工作。供
|
|
* useTaskProgressLoop 的可选协调选项复用。
|
|
*/
|
|
export interface TaskPollingCoordinatorOptions {
|
|
/** 任务集合最大容量,默认 500 */
|
|
maxTasks?: number
|
|
/** 终态清理回调;同一 taskId 只回调一次 */
|
|
onTerminal?: (taskId: number) => void
|
|
}
|
|
|
|
export interface TaskPollingCoordinatorStats {
|
|
/** 成功加入的任务数 */
|
|
addCount: number
|
|
/** 因重复被去重掉的 add 数 */
|
|
dedupedCount: number
|
|
/** 因超过 maxTasks 被拒绝的 add 数 */
|
|
rejectedCount: number
|
|
/** 触发终态清理的任务数 */
|
|
terminalCount: number
|
|
/** 实际发出的请求次数 */
|
|
requestCount: number
|
|
/** 被合并进 in-flight 请求的 runOnce 调用数 */
|
|
mergedCount: number
|
|
}
|
|
|
|
export function createTaskPollingCoordinator(options: TaskPollingCoordinatorOptions = {}) {
|
|
const maxTasks = options.maxTasks ?? 500
|
|
const onTerminal = options.onTerminal ?? (() => {})
|
|
if (!(maxTasks > 0)) {
|
|
throw new Error('maxTasks 必须为正数: ' + maxTasks)
|
|
}
|
|
|
|
const set = new Set<number>()
|
|
let addCount = 0
|
|
let dedupedCount = 0
|
|
let rejectedCount = 0
|
|
let terminalCount = 0
|
|
let requestCount = 0
|
|
let mergedCount = 0
|
|
let inFlight: Promise<unknown> | null = null
|
|
|
|
function validateTaskId(taskId: number): boolean {
|
|
return Number.isFinite(taskId) && taskId > 0 && Number.isInteger(taskId)
|
|
}
|
|
|
|
function add(taskId: number): boolean {
|
|
if (!validateTaskId(taskId)) {
|
|
throw new Error('taskId 必须是正整数: ' + taskId)
|
|
}
|
|
if (set.has(taskId)) {
|
|
dedupedCount += 1
|
|
return false
|
|
}
|
|
if (set.size >= maxTasks) {
|
|
rejectedCount += 1
|
|
return false
|
|
}
|
|
set.add(taskId)
|
|
addCount += 1
|
|
return true
|
|
}
|
|
|
|
function remove(taskId: number): boolean {
|
|
if (!validateTaskId(taskId)) return false
|
|
return set.delete(taskId)
|
|
}
|
|
|
|
function markTerminal(taskId: number) {
|
|
if (!validateTaskId(taskId)) return
|
|
if (set.delete(taskId)) {
|
|
terminalCount += 1
|
|
onTerminal(taskId)
|
|
}
|
|
}
|
|
|
|
function clear() {
|
|
set.clear()
|
|
}
|
|
|
|
function has(taskId: number): boolean {
|
|
return set.has(taskId)
|
|
}
|
|
|
|
function ids(): number[] {
|
|
return [...set]
|
|
}
|
|
|
|
function runOnce<T>(loader: (taskIds: number[]) => Promise<T>): Promise<T> {
|
|
if (inFlight) {
|
|
mergedCount += 1
|
|
return inFlight as Promise<T>
|
|
}
|
|
const taskIds = ids()
|
|
if (taskIds.length === 0) {
|
|
return Promise.resolve(undefined as unknown as T)
|
|
}
|
|
requestCount += 1
|
|
inFlight = loader(taskIds).finally(() => {
|
|
inFlight = null
|
|
})
|
|
return inFlight as Promise<T>
|
|
}
|
|
|
|
function stats(): TaskPollingCoordinatorStats {
|
|
return {
|
|
addCount,
|
|
dedupedCount,
|
|
rejectedCount,
|
|
terminalCount,
|
|
requestCount,
|
|
mergedCount,
|
|
}
|
|
}
|
|
|
|
return {
|
|
get size() {
|
|
return set.size
|
|
},
|
|
add,
|
|
remove,
|
|
markTerminal,
|
|
clear,
|
|
has,
|
|
ids,
|
|
runOnce,
|
|
stats,
|
|
}
|
|
}
|
|
|
|
export type TaskPollingCoordinator = ReturnType<typeof createTaskPollingCoordinator>
|