Files
crawler-plugin/frontend-vue/src/shared/composables/useTaskProgressLoop.ts
T
huangzd1997 58e91589df task-86: 将隐藏页面轮询间隔、前台恢复和退避策略统一配置化
task-progress-config 新增统一配置 API:getTaskPollBackoffMs(attempt)
指数退避(base×2^attempt 封顶 max)、前台恢复开关与延迟
(getTaskForegroundRefreshEnabled/DelayMs)、configureTaskPolling/
resetTaskPollingConfig 运行时配置。校验失败抛错且零状态变更,
document 缺失时按可见间隔降级。useTaskProgressLoop 改用退避配置
替代硬编码 500ms 重试,前台恢复走可配置开关与延迟(含 in-flight 防
并发)。8 个测试覆盖默认、多配置、幂等、空、单字段、封顶、非法输入
与故障恢复。
2026-08-30 22:51:21 +08:00

355 lines
11 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,
getTaskPollBackoffMs,
getTaskForegroundRefreshEnabled,
getTaskForegroundRefreshDelayMs,
} 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, getTaskPollBackoffMs(0))
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 (
getTaskForegroundRefreshEnabled() &&
document.visibilityState === 'visible' &&
taskIds.value.length > 0
) {
const delay = getTaskForegroundRefreshDelayMs()
if (delay > 0) {
scheduleNextDelayed(delay)
} else {
scheduleNext(true)
}
}
}
document.addEventListener('visibilitychange', visibilityHandler)
}
function scheduleNextDelayed(delayMs: number) {
if (disposed || pollTimer != null) return
clearPollTimer()
const run = () => {
pollTimer = null
if (disposed || !taskIds.value.length) return
if (inFlight.value) {
pollTimer = timers.setTimeout('task-poll', run, getTaskPollBackoffMs(0))
return
}
void refreshOnce()
if (!disposed && taskIds.value.length > 0) {
pollTimer = timers.setTimeout('task-poll', run, intervalMs())
}
}
pollTimer = timers.setTimeout('task-poll', run, delayMs)
}
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,
}
}