e60bd8567a
新增纯 TS 任务轮询协调器 task-polling-coordinator:重复 add 去重、 并发 runOnce 合并为单次请求并共享结果、markTerminal 幂等终态清理 (只回调一次)、maxTasks 上限拒绝超出条目。useTaskProgressLoop 新增 可选 coordinator 选项,add/remove/终态路径委托协调器统一处理(默认 不启用)。测试用 node:test 覆盖 8 个场景,含 3 路并发合并断言。
188 lines
6.7 KiB
TypeScript
188 lines
6.7 KiB
TypeScript
import { test } from 'node:test'
|
|
import assert from 'node:assert/strict'
|
|
import { createTaskPollingCoordinator } from '../src/shared/task-polling-coordinator.ts'
|
|
|
|
function deferred() {
|
|
let resolve!: (v: unknown) => void
|
|
let reject!: (e: unknown) => void
|
|
const promise = new Promise((res, rej) => {
|
|
resolve = res
|
|
reject = rej
|
|
})
|
|
return { promise, resolve, reject }
|
|
}
|
|
|
|
test('test_task_083_merge_cleanup_polling_normal_default_path', async () => {
|
|
const terminal: number[] = []
|
|
const coordinator = createTaskPollingCoordinator({ onTerminal: (id) => terminal.push(id) })
|
|
assert.equal(coordinator.add(1), true)
|
|
assert.equal(coordinator.add(2), true)
|
|
assert.equal(coordinator.size, 2)
|
|
assert.deepEqual(coordinator.ids(), [1, 2])
|
|
const calls: number[][] = []
|
|
const result = await coordinator.runOnce((ids) => {
|
|
calls.push(ids)
|
|
return Promise.resolve({ ok: true, n: ids.length })
|
|
})
|
|
assert.deepEqual(result, { ok: true, n: 2 })
|
|
assert.equal(calls.length, 1)
|
|
assert.deepEqual(calls[0], [1, 2])
|
|
coordinator.markTerminal(1)
|
|
assert.equal(coordinator.has(1), false)
|
|
assert.deepEqual(coordinator.ids(), [2])
|
|
assert.deepEqual(terminal, [1])
|
|
const stats = coordinator.stats()
|
|
assert.equal(stats.addCount, 2)
|
|
assert.equal(stats.dedupedCount, 0)
|
|
assert.equal(stats.terminalCount, 1)
|
|
assert.equal(stats.requestCount, 1)
|
|
assert.equal(stats.mergedCount, 0)
|
|
assert.equal(stats.rejectedCount, 0)
|
|
})
|
|
|
|
test('test_task_083_merge_cleanup_polling_normal_multiple_items', async () => {
|
|
const terminal: number[] = []
|
|
const coordinator = createTaskPollingCoordinator({ onTerminal: (id) => terminal.push(id) })
|
|
const ids = Array.from({ length: 10 }, (_, i) => i + 1)
|
|
const added = ids.map((id) => coordinator.add(id))
|
|
assert.deepEqual(added, Array(10).fill(true))
|
|
assert.equal(coordinator.size, 10)
|
|
assert.deepEqual(coordinator.ids(), ids)
|
|
const seen: number[][] = []
|
|
await coordinator.runOnce((requested) => {
|
|
seen.push(requested)
|
|
return Promise.resolve(undefined)
|
|
})
|
|
assert.deepEqual(seen, [ids])
|
|
for (const id of ids) coordinator.markTerminal(id)
|
|
assert.equal(coordinator.size, 0)
|
|
assert.deepEqual(terminal, ids)
|
|
assert.equal(coordinator.stats().terminalCount, 10)
|
|
})
|
|
|
|
test('test_task_083_merge_cleanup_polling_normal_repeated_operation_is_idempotent', async () => {
|
|
const terminal: number[] = []
|
|
const coordinator = createTaskPollingCoordinator({ onTerminal: (id) => terminal.push(id) })
|
|
coordinator.add(1)
|
|
assert.equal(coordinator.add(1), false)
|
|
assert.equal(coordinator.add(1), false)
|
|
assert.equal(coordinator.size, 1)
|
|
assert.deepEqual(coordinator.ids(), [1])
|
|
assert.equal(coordinator.stats().dedupedCount, 2)
|
|
// 重复终态清理只回调一次
|
|
coordinator.markTerminal(1)
|
|
coordinator.markTerminal(1)
|
|
coordinator.markTerminal(1)
|
|
assert.deepEqual(terminal, [1])
|
|
assert.equal(coordinator.stats().terminalCount, 1)
|
|
// 并发 runOnce 合并为单次请求
|
|
coordinator.add(1)
|
|
const calls: number[][] = []
|
|
const d = deferred()
|
|
const p1 = coordinator.runOnce((ids) => {
|
|
calls.push(ids)
|
|
return d.promise
|
|
})
|
|
const p2 = coordinator.runOnce((ids) => {
|
|
calls.push(ids)
|
|
return Promise.resolve('second')
|
|
})
|
|
const p3 = coordinator.runOnce((ids) => {
|
|
calls.push(ids)
|
|
return Promise.resolve('third')
|
|
})
|
|
d.resolve('first')
|
|
const results = await Promise.all([p1, p2, p3])
|
|
assert.deepEqual(results, ['first', 'first', 'first'])
|
|
assert.equal(calls.length, 1)
|
|
const stats = coordinator.stats()
|
|
assert.equal(stats.requestCount, 1)
|
|
assert.equal(stats.mergedCount, 2)
|
|
})
|
|
|
|
test('test_task_083_merge_cleanup_polling_boundary_empty_input', async () => {
|
|
const coordinator = createTaskPollingCoordinator()
|
|
assert.equal(coordinator.size, 0)
|
|
assert.deepEqual(coordinator.ids(), [])
|
|
let called = false
|
|
const result = await coordinator.runOnce(() => {
|
|
called = true
|
|
return Promise.resolve({ n: 1 })
|
|
})
|
|
assert.equal(result, undefined)
|
|
assert.equal(called, false)
|
|
assert.equal(coordinator.stats().requestCount, 0)
|
|
coordinator.clear()
|
|
coordinator.markTerminal(999)
|
|
assert.equal(coordinator.stats().terminalCount, 0)
|
|
})
|
|
|
|
test('test_task_083_merge_cleanup_polling_boundary_single_item', async () => {
|
|
const coordinator = createTaskPollingCoordinator()
|
|
assert.equal(coordinator.add(7), true)
|
|
const calls: number[][] = []
|
|
await coordinator.runOnce((ids) => {
|
|
calls.push(ids)
|
|
return Promise.resolve(undefined)
|
|
})
|
|
assert.deepEqual(calls, [[7]])
|
|
assert.equal(coordinator.size, 1)
|
|
assert.equal(coordinator.stats().requestCount, 1)
|
|
})
|
|
|
|
test('test_task_083_merge_cleanup_polling_boundary_limit_and_overflow', () => {
|
|
const coordinator = createTaskPollingCoordinator({ maxTasks: 3 })
|
|
assert.equal(coordinator.add(1), true)
|
|
assert.equal(coordinator.add(2), true)
|
|
assert.equal(coordinator.add(3), true)
|
|
assert.equal(coordinator.size, 3)
|
|
assert.equal(coordinator.add(4), false)
|
|
assert.equal(coordinator.add(5), false)
|
|
assert.equal(coordinator.size, 3)
|
|
assert.deepEqual(coordinator.ids(), [1, 2, 3])
|
|
const stats = coordinator.stats()
|
|
assert.equal(stats.rejectedCount, 2)
|
|
assert.equal(stats.addCount, 3)
|
|
// 移除后可继续加入
|
|
coordinator.remove(1)
|
|
assert.equal(coordinator.add(4), true)
|
|
assert.deepEqual(coordinator.ids(), [2, 3, 4])
|
|
})
|
|
|
|
test('test_task_083_merge_cleanup_polling_invalid_input_rejected', () => {
|
|
const coordinator = createTaskPollingCoordinator()
|
|
assert.throws(() => coordinator.add(0), /taskId 必须是正整数/)
|
|
assert.throws(() => coordinator.add(-1), /taskId 必须是正整数/)
|
|
assert.throws(() => coordinator.add(NaN), /taskId 必须是正整数/)
|
|
assert.throws(() => coordinator.add(1.5), /taskId 必须是正整数/)
|
|
assert.equal(coordinator.size, 0)
|
|
assert.equal(coordinator.stats().rejectedCount, 0)
|
|
// 读取路径宽容
|
|
assert.equal(coordinator.has(0), false)
|
|
assert.equal(coordinator.remove(0), false)
|
|
assert.throws(() => createTaskPollingCoordinator({ maxTasks: 0 }), /maxTasks 必须为正数/)
|
|
})
|
|
|
|
test('test_task_083_merge_cleanup_polling_dependency_failure_releases_resources', async () => {
|
|
const coordinator = createTaskPollingCoordinator()
|
|
coordinator.add(1)
|
|
coordinator.add(2)
|
|
// loader 抛错:请求失败但任务集合保持不变,in-flight 释放
|
|
await assert.rejects(
|
|
coordinator.runOnce(() => Promise.reject(new Error('network down'))),
|
|
/network down/,
|
|
)
|
|
assert.equal(coordinator.size, 2)
|
|
// 恢复后同一实例可再次发起请求
|
|
const calls: number[][] = []
|
|
const result = await coordinator.runOnce((ids) => {
|
|
calls.push(ids)
|
|
return Promise.resolve('recovered')
|
|
})
|
|
assert.equal(result, 'recovered')
|
|
assert.deepEqual(calls, [[1, 2]])
|
|
const stats = coordinator.stats()
|
|
assert.equal(stats.requestCount, 2)
|
|
assert.equal(stats.mergedCount, 0)
|
|
})
|