Files
crawler-plugin/frontend-vue/src/shared/composables/useTaskProgressLoop.ts
T
huangzd1997 e60bd8567a task-83: 统一不同页面的轮询去重、in-flight 合并和终态清理
新增纯 TS 任务轮询协调器 task-polling-coordinator:重复 add 去重、
并发 runOnce 合并为单次请求并共享结果、markTerminal 幂等终态清理
(只回调一次)、maxTasks 上限拒绝超出条目。useTaskProgressLoop 新增
可选 coordinator 选项,add/remove/终态路径委托协调器统一处理(默认
不启用)。测试用 node:test 覆盖 8 个场景,含 3 路并发合并断言。
2026-08-30 22:44:12 +08:00

323 lines
10 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import { onBeforeUnmount, ref, watch, type Ref } from 'vue'
import { getTaskPollIntervalMs } from '@/shared/task-progress-config'
import { createCategorizedTimers } from '@/shared/utils/categorized-timers'
import {
createTaskPollingBaseline,
type TaskPollingBaseline,
} from '@/shared/task-polling-baseline'
import {
createProgressResponseCache,
type ProgressResponseCache,
} from '@/shared/progress-response-cache'
import {
createTaskPollingCoordinator,
type TaskPollingCoordinator,
} from '@/shared/task-polling-coordinator'
/**
* 通用任务进度轮询组合式函数。
*
* 各模块统一使用:保存 taskId 列表(可选 localStorage 持久化),按当前可见性
* 周期性请求批量进度接口,每个任务到达终态时回调上层做状态同步、列表刷新等。
*
* 使用方式:
* const loop = useTaskProgressLoop<MyTaskDetailVo>({
* scope: 'collect-data',
* storageKey: 'brand:collect-data:polling-task-ids',
* fetchProgress: (ids) => getCollectDataTaskProgressBatch(ids),
* extractTaskId: (detail) => detail.task?.id,
* extractStatus: (detail) => detail.task?.status,
* isTerminal: (status) => status === 'SUCCESS' || status === 'FAILED',
* onUpdate: (taskId, detail) => { ... },
* onTerminal: async (taskId, detail) => { await refreshHistory() },
* })
* loop.add(taskId) // 触发轮询
* loop.dispose() // 组件卸载时调用(自动通过 onBeforeUnmount 清理)
*/
export interface TaskProgressLoopOptions<TDetail> {
/** 用于 categorized-timers 的命名空间,例如 'collect-data';同一页面内须唯一 */
scope: string
/** localStorage 持久化的 key;省略则不持久化 */
storageKey?: string
/** 拉取批量进度的接口;返回 items 数组 */
fetchProgress: (taskIds: number[]) => Promise<{ items?: TDetail[] }>
/** 从单条进度详情中提取 taskId */
extractTaskId: (detail: TDetail) => number | null | undefined
/** 从单条进度详情中提取状态字符串(如 'PENDING'/'RUNNING'/'SUCCESS'/'FAILED' */
extractStatus: (detail: TDetail) => string | null | undefined
/** 判定是否终态;默认 SUCCESS / FAILED 视为终态 */
isTerminal?: (status: string) => boolean
/** 每条进度落地时调用,用于上层缓存最新快照 */
onUpdate?: (taskId: number, detail: TDetail) => void
/**
* 任意一个任务进入终态时调用;可以是 async(轮询会等待完成再调度下一轮)。
* 若多个任务同一轮到达终态,会被分别回调。
*/
onTerminal?: (taskId: number, detail: TDetail | undefined, status: string) => void | Promise<void>
/** 轮询周期失败时的回调;默认静默 */
onError?: (error: unknown) => void
/** 自定义轮询间隔;默认根据 document.visibilityState 自适应(5s/30s */
getIntervalMs?: () => number
/**
* 可选轮询基线统计(Task 81):传入后每轮请求与响应都会被记录到基线,
* 并在每次 refreshOnce 结束后把基线实例回传,供页面/测试观测请求量、
* 响应体大小与缓存内存占用。
*/
onBaseline?: (baseline: TaskPollingBaseline) => void
/**
* 可选进度响应缓存(Task 82):传入后每轮响应快照按 taskId 写入
* TTL+最大条目数的有界缓存(读取时惰性清理过期条目),供页面在
* 轮询间隙复用最近一次进度快照。
*/
progressCache?: ProgressResponseCache<TDetail>
/**
* 可选轮询协调器(Task 83):传入后 add/remove/终态清理委托给协调器
* 统一去重、合并并发请求;终态任务经协调器回调后从轮询集合移除,
* 同一任务只触发一次终态回调。
*/
coordinator?: TaskPollingCoordinator
}
export interface TaskProgressLoopHandle<TDetail> {
taskIds: Ref<number[]>
taskStatuses: Ref<Record<number, string>>
inFlight: Ref<boolean>
add: (taskId: number) => void
remove: (taskId: number) => void
reset: (taskIds: number[]) => void
ensure: (immediate?: boolean) => void
stop: () => void
refreshOnce: () => Promise<void>
isTerminal: (taskId: number) => boolean
dispose: () => void
}
const DEFAULT_TERMINAL = (status: string) => status === 'SUCCESS' || status === 'FAILED'
function readIdsFromStorage(key?: string): number[] {
if (!key || typeof window === 'undefined') return []
try {
const raw = window.localStorage.getItem(key)
if (!raw) return []
const parsed = JSON.parse(raw) as unknown
return Array.isArray(parsed) ? parsed.filter((n): n is number => typeof n === 'number' && n > 0) : []
} catch {
return []
}
}
function writeIdsToStorage(key: string | undefined, ids: number[]) {
if (!key || typeof window === 'undefined') return
try {
if (ids.length === 0) {
window.localStorage.removeItem(key)
} else {
window.localStorage.setItem(key, JSON.stringify(ids))
}
} catch {
/* 写本地存储失败不影响功能 */
}
}
export function useTaskProgressLoop<TDetail>(
options: TaskProgressLoopOptions<TDetail>,
): TaskProgressLoopHandle<TDetail> {
const timers = createCategorizedTimers(`task-progress-loop:${options.scope}`)
const isTerminal = options.isTerminal ?? DEFAULT_TERMINAL
const intervalMs = options.getIntervalMs ?? getTaskPollIntervalMs
const baseline = options.onBaseline ? createTaskPollingBaseline() : null
const coordinator = options.coordinator ?? null
const taskIds = ref<number[]>(readIdsFromStorage(options.storageKey))
const taskStatuses = ref<Record<number, string>>({})
const inFlight = ref(false)
let pollTimer: number | null = null
let disposed = false
function persist() {
writeIdsToStorage(options.storageKey, taskIds.value)
}
function add(taskId: number) {
if (!Number.isFinite(taskId) || taskId <= 0) return
if (coordinator && !coordinator.add(taskId)) return
if (taskIds.value.includes(taskId)) return
taskIds.value = [...taskIds.value, taskId]
persist()
ensure(true)
}
function remove(taskId: number) {
coordinator?.remove(taskId)
if (!taskIds.value.includes(taskId)) return
taskIds.value = taskIds.value.filter((id) => id !== taskId)
persist()
if (taskStatuses.value[taskId]) {
const next = { ...taskStatuses.value }
delete next[taskId]
taskStatuses.value = next
}
}
function reset(ids: number[]) {
const cleaned = Array.from(new Set(ids.filter((n) => Number.isFinite(n) && n > 0)))
taskIds.value = cleaned
persist()
}
function isTerminalById(taskId: number) {
const s = taskStatuses.value[taskId]
return !!s && isTerminal(s)
}
async function refreshOnce() {
if (disposed) return
const ids = taskIds.value.filter((id) => id > 0)
if (!ids.length) return
inFlight.value = true
baseline?.recordRequest()
try {
const result = await options.fetchProgress(ids)
const items = result?.items || []
const terminalEvents: Array<{ taskId: number; detail: TDetail | undefined; status: string }> = []
const nextStatuses = { ...taskStatuses.value }
for (const detail of items) {
const id = options.extractTaskId(detail)
if (typeof id !== 'number' || id <= 0) continue
const status = options.extractStatus(detail) || ''
try {
options.onUpdate?.(id, detail)
} catch {
/* onUpdate 抛错不应中断本轮 */
}
options.progressCache?.set(id, detail)
if (status) {
nextStatuses[id] = status
if (isTerminal(status)) {
terminalEvents.push({ taskId: id, detail, status })
}
}
}
taskStatuses.value = nextStatuses
for (const event of terminalEvents) {
coordinator?.markTerminal(event.taskId)
remove(event.taskId)
try {
await options.onTerminal?.(event.taskId, event.detail, event.status)
} catch {
/* onTerminal 抛错只影响一次回调 */
}
}
} catch (error) {
options.onError?.(error)
} finally {
inFlight.value = false
if (baseline) options.onBaseline?.(baseline)
}
}
function clearPollTimer() {
if (pollTimer != null) {
timers.clearTimer('task-poll', pollTimer)
pollTimer = null
}
}
function scheduleNext(immediate = false) {
if (disposed) return
if (pollTimer != null && !immediate) return
clearPollTimer()
const run = async () => {
pollTimer = null
if (disposed) return
if (!taskIds.value.length) return
if (inFlight.value) {
// 上一次还没回,500ms 后再试
pollTimer = timers.setTimeout('task-poll', run, 500)
return
}
await refreshOnce()
if (!disposed && taskIds.value.length > 0) {
pollTimer = timers.setTimeout('task-poll', run, intervalMs())
}
}
if (immediate) {
void run()
} else {
pollTimer = timers.setTimeout('task-poll', run, intervalMs())
}
}
function ensure(immediate = false) {
if (disposed) return
if (!taskIds.value.length) return
if (pollTimer != null && !immediate) return
scheduleNext(immediate)
}
function stop() {
clearPollTimer()
}
// 任务列表清空时自动停止;新增时自动启动一轮
watch(
taskIds,
(ids, prev) => {
if (disposed) return
if (!ids.length) {
stop()
return
}
if (!prev || prev.length === 0) {
scheduleNext(true)
}
},
{ flush: 'post' },
)
// 切到前台后立刻拉一次,让用户回到页面看到的是最新状态
let visibilityHandler: (() => void) | null = null
if (typeof document !== 'undefined') {
visibilityHandler = () => {
if (document.visibilityState === 'visible' && taskIds.value.length > 0) {
scheduleNext(true)
}
}
document.addEventListener('visibilitychange', visibilityHandler)
}
function dispose() {
if (disposed) return
disposed = true
stop()
timers.clearScope()
if (visibilityHandler && typeof document !== 'undefined') {
document.removeEventListener('visibilitychange', visibilityHandler)
visibilityHandler = null
}
}
onBeforeUnmount(() => dispose())
// 初始化时若已有任务(从 storage 恢复),立即开始一轮
if (taskIds.value.length > 0) {
scheduleNext(true)
}
return {
taskIds,
taskStatuses,
inFlight,
add,
remove,
reset,
ensure,
stop,
refreshOnce,
isTerminal: isTerminalById,
dispose,
}
}