[backend] fix: 修复异常处理和类型转换问题

This commit is contained in:
xsl
2026-01-26 11:53:40 +08:00
parent 7ccc2a6ac6
commit 83e05bf85f
28639 changed files with 2506458 additions and 93 deletions
+310
View File
@@ -0,0 +1,310 @@
# fastq
![ci][ci-url]
[![npm version][npm-badge]][npm-url]
Fast, in memory work queue.
Benchmarks (1 million tasks):
* setImmediate: 812ms
* fastq: 854ms
* async.queue: 1298ms
* neoAsync.queue: 1249ms
Obtained on node 12.16.1, on a dedicated server.
If you need zero-overhead series function call, check out
[fastseries](http://npm.im/fastseries). For zero-overhead parallel
function call, check out [fastparallel](http://npm.im/fastparallel).
* <a href="#install">Installation</a>
* <a href="#usage">Usage</a>
* <a href="#api">API</a>
* <a href="#license">Licence &amp; copyright</a>
## Install
`npm i fastq --save`
## Usage (callback API)
```js
'use strict'
const queue = require('fastq')(worker, 1)
queue.push(42, function (err, result) {
if (err) { throw err }
console.log('the result is', result)
})
function worker (arg, cb) {
cb(null, arg * 2)
}
```
## Usage (promise API)
```js
const queue = require('fastq').promise(worker, 1)
async function worker (arg) {
return arg * 2
}
async function run () {
const result = await queue.push(42)
console.log('the result is', result)
}
run()
```
### Setting "this"
```js
'use strict'
const that = { hello: 'world' }
const queue = require('fastq')(that, worker, 1)
queue.push(42, function (err, result) {
if (err) { throw err }
console.log(this)
console.log('the result is', result)
})
function worker (arg, cb) {
console.log(this)
cb(null, arg * 2)
}
```
### Using with TypeScript (callback API)
```ts
'use strict'
import * as fastq from "fastq";
import type { queue, done } from "fastq";
type Task = {
id: number
}
const q: queue<Task> = fastq(worker, 1)
q.push({ id: 42})
function worker (arg: Task, cb: done) {
console.log(arg.id)
cb(null)
}
```
### Using with TypeScript (promise API)
```ts
'use strict'
import * as fastq from "fastq";
import type { queueAsPromised } from "fastq";
type Task = {
id: number
}
const q: queueAsPromised<Task> = fastq.promise(asyncWorker, 1)
q.push({ id: 42}).catch((err) => console.error(err))
async function asyncWorker (arg: Task): Promise<void> {
// No need for a try-catch block, fastq handles errors automatically
console.log(arg.id)
}
```
## API
* <a href="#fastqueue"><code>fastqueue()</code></a>
* <a href="#push"><code>queue#<b>push()</b></code></a>
* <a href="#unshift"><code>queue#<b>unshift()</b></code></a>
* <a href="#pause"><code>queue#<b>pause()</b></code></a>
* <a href="#resume"><code>queue#<b>resume()</b></code></a>
* <a href="#idle"><code>queue#<b>idle()</b></code></a>
* <a href="#length"><code>queue#<b>length()</b></code></a>
* <a href="#getQueue"><code>queue#<b>getQueue()</b></code></a>
* <a href="#kill"><code>queue#<b>kill()</b></code></a>
* <a href="#killAndDrain"><code>queue#<b>killAndDrain()</b></code></a>
* <a href="#error"><code>queue#<b>error()</b></code></a>
* <a href="#concurrency"><code>queue#<b>concurrency</b></code></a>
* <a href="#drain"><code>queue#<b>drain</b></code></a>
* <a href="#empty"><code>queue#<b>empty</b></code></a>
* <a href="#saturated"><code>queue#<b>saturated</b></code></a>
* <a href="#promise"><code>fastqueue.promise()</code></a>
-------------------------------------------------------
<a name="fastqueue"></a>
### fastqueue([that], worker, concurrency)
Creates a new queue.
Arguments:
* `that`, optional context of the `worker` function.
* `worker`, worker function, it would be called with `that` as `this`,
if that is specified.
* `concurrency`, number of concurrent tasks that could be executed in
parallel.
-------------------------------------------------------
<a name="push"></a>
### queue.push(task, done)
Add a task at the end of the queue. `done(err, result)` will be called
when the task was processed.
-------------------------------------------------------
<a name="unshift"></a>
### queue.unshift(task, done)
Add a task at the beginning of the queue. `done(err, result)` will be called
when the task was processed.
-------------------------------------------------------
<a name="pause"></a>
### queue.pause()
Pause the processing of tasks. Currently worked tasks are not
stopped.
-------------------------------------------------------
<a name="resume"></a>
### queue.resume()
Resume the processing of tasks.
-------------------------------------------------------
<a name="idle"></a>
### queue.idle()
Returns `false` if there are tasks being processed or waiting to be processed.
`true` otherwise.
-------------------------------------------------------
<a name="length"></a>
### queue.length()
Returns the number of tasks waiting to be processed (in the queue).
-------------------------------------------------------
<a name="getQueue"></a>
### queue.getQueue()
Returns all the tasks be processed (in the queue). Returns empty array when there are no tasks
-------------------------------------------------------
<a name="kill"></a>
### queue.kill()
Removes all tasks waiting to be processed, and reset `drain` to an empty
function.
-------------------------------------------------------
<a name="killAndDrain"></a>
### queue.killAndDrain()
Same than `kill` but the `drain` function will be called before reset to empty.
-------------------------------------------------------
<a name="error"></a>
### queue.error(handler)
Set a global error handler. `handler(err, task)` will be called
each time a task is completed, `err` will be not null if the task has thrown an error.
-------------------------------------------------------
<a name="concurrency"></a>
### queue.concurrency
Property that returns the number of concurrent tasks that could be executed in
parallel. It can be altered at runtime.
-------------------------------------------------------
<a name="paused"></a>
### queue.paused
Property (Read-Only) that returns `true` when the queue is in a paused state.
-------------------------------------------------------
<a name="drain"></a>
### queue.drain
Function that will be called when the last
item from the queue has been processed by a worker.
It can be altered at runtime.
-------------------------------------------------------
<a name="empty"></a>
### queue.empty
Function that will be called when the last
item from the queue has been assigned to a worker.
It can be altered at runtime.
-------------------------------------------------------
<a name="saturated"></a>
### queue.saturated
Function that will be called when the queue hits the concurrency
limit.
It can be altered at runtime.
-------------------------------------------------------
<a name="promise"></a>
### fastqueue.promise([that], worker(arg), concurrency)
Creates a new queue with `Promise` apis. It also offers all the methods
and properties of the object returned by [`fastqueue`](#fastqueue) with the modified
[`push`](#pushPromise) and [`unshift`](#unshiftPromise) methods.
Node v10+ is required to use the promisified version.
Arguments:
* `that`, optional context of the `worker` function.
* `worker`, worker function, it would be called with `that` as `this`,
if that is specified. It MUST return a `Promise`.
* `concurrency`, number of concurrent tasks that could be executed in
parallel.
<a name="pushPromise"></a>
#### queue.push(task) => Promise
Add a task at the end of the queue. The returned `Promise` will be fulfilled (rejected)
when the task is completed successfully (unsuccessfully).
This promise could be ignored as it will not lead to a `'unhandledRejection'`.
<a name="unshiftPromise"></a>
#### queue.unshift(task) => Promise
Add a task at the beginning of the queue. The returned `Promise` will be fulfilled (rejected)
when the task is completed successfully (unsuccessfully).
This promise could be ignored as it will not lead to a `'unhandledRejection'`.
<a name="drained"></a>
#### queue.drained() => Promise
Wait for the queue to be drained. The returned `Promise` will be resolved when all tasks in the queue have been processed by a worker.
This promise could be ignored as it will not lead to a `'unhandledRejection'`.
## License
ISC
[ci-url]: https://github.com/mcollina/fastq/workflows/ci/badge.svg
[npm-badge]: https://badge.fury.io/js/fastq.svg
[npm-url]: https://badge.fury.io/js/fastq
+15
View File
@@ -0,0 +1,15 @@
# Security Policy
## Supported Versions
Use this section to tell people about which versions of your project are
currently being supported with security updates.
| Version | Supported |
| ------- | ------------------ |
| 1.x | :white_check_mark: |
| < 1.0 | :x: |
## Reporting a Vulnerability
Please report all vulnerabilities at [https://github.com/mcollina/fastq/security](https://github.com/mcollina/fastq/security).
+9
View File
@@ -0,0 +1,9 @@
import { promise as queueAsPromised } from './queue.js'
const queue = queueAsPromised(worker, 1)
console.log('the result is', await queue.push(42))
async function worker (arg) {
return 42 * 2
}
+59
View File
@@ -0,0 +1,59 @@
declare function fastq<C, T = any, R = any>(context: C, worker: fastq.worker<C, T, R>, concurrency: number): fastq.queue<T, R>
declare function fastq<C, T = any, R = any>(worker: fastq.worker<C, T, R>, concurrency: number): fastq.queue<T, R>
declare namespace fastq {
type worker<C, T = any, R = any> = (this: C, task: T, cb: fastq.done<R>) => void
type asyncWorker<C, T = any, R = any> = (this: C, task: T) => Promise<R>
type done<R = any> = (err: Error | null, result?: R) => void
type errorHandler<T = any> = (err: Error, task: T) => void
interface queue<T = any, R = any> {
/** Add a task at the end of the queue. `done(err, result)` will be called when the task was processed. */
push(task: T, done?: done<R>): void
/** Add a task at the beginning of the queue. `done(err, result)` will be called when the task was processed. */
unshift(task: T, done?: done<R>): void
/** Pause the processing of tasks. Currently worked tasks are not stopped. */
pause(): any
/** Resume the processing of tasks. */
resume(): any
running(): number
/** Returns `false` if there are tasks being processed or waiting to be processed. `true` otherwise. */
idle(): boolean
/** Returns the number of tasks waiting to be processed (in the queue). */
length(): number
/** Returns all the tasks be processed (in the queue). Returns empty array when there are no tasks */
getQueue(): T[]
/** Removes all tasks waiting to be processed, and reset `drain` to an empty function. */
kill(): any
/** Same than `kill` but the `drain` function will be called before reset to empty. */
killAndDrain(): any
/** Removes all tasks waiting to be processed, calls each task's callback with an abort error (rejects promises for promise-based queues), and resets `drain` to an empty function. */
abort(): any
/** Set a global error handler. `handler(err, task)` will be called each time a task is completed, `err` will be not null if the task has thrown an error. */
error(handler: errorHandler<T>): void
/** Property that returns the number of concurrent tasks that could be executed in parallel. It can be altered at runtime. */
concurrency: number
/** Property (Read-Only) that returns `true` when the queue is in a paused state. */
readonly paused: boolean
/** Function that will be called when the last item from the queue has been processed by a worker. It can be altered at runtime. */
drain(): any
/** Function that will be called when the last item from the queue has been assigned to a worker. It can be altered at runtime. */
empty: () => void
/** Function that will be called when the queue hits the concurrency limit. It can be altered at runtime. */
saturated: () => void
}
interface queueAsPromised<T = any, R = any> extends queue<T, R> {
/** Add a task at the end of the queue. The returned `Promise` will be fulfilled (rejected) when the task is completed successfully (unsuccessfully). */
push(task: T): Promise<R>
/** Add a task at the beginning of the queue. The returned `Promise` will be fulfilled (rejected) when the task is completed successfully (unsuccessfully). */
unshift(task: T): Promise<R>
/** Wait for the queue to be drained. The returned `Promise` will be resolved when all tasks in the queue have been processed by a worker. */
drained(): Promise<void>
}
function promise<C, T = any, R = any>(context: C, worker: fastq.asyncWorker<C, T, R>, concurrency: number): fastq.queueAsPromised<T, R>
function promise<C, T = any, R = any>(worker: fastq.asyncWorker<C, T, R>, concurrency: number): fastq.queueAsPromised<T, R>
}
export = fastq
+49
View File
@@ -0,0 +1,49 @@
{
"name": "fastq",
"version": "1.20.1",
"description": "Fast, in memory work queue",
"main": "queue.js",
"type": "commonjs",
"scripts": {
"lint": "eslint .",
"unit": "nyc --lines 100 --branches 100 --functions 100 --check-coverage --reporter=text tape test/test.js test/promise.js",
"coverage": "nyc --reporter=html --reporter=cobertura --reporter=text tape test/test.js test/promise.js",
"test:report": "npm run lint && npm run unit:report",
"test": "npm run lint && npm run unit",
"typescript": "tsc --project ./test/tsconfig.json",
"legacy": "tape test/test.js"
},
"pre-commit": [
"test",
"typescript"
],
"repository": {
"type": "git",
"url": "git+https://github.com/mcollina/fastq.git"
},
"keywords": [
"fast",
"queue",
"async",
"worker"
],
"author": "Matteo Collina <hello@matteocollina.com>",
"license": "ISC",
"bugs": {
"url": "https://github.com/mcollina/fastq/issues"
},
"homepage": "https://github.com/mcollina/fastq#readme",
"devDependencies": {
"async": "^3.1.0",
"eslint": "^9.36.0",
"neo-async": "^2.6.1",
"neostandard": "^0.12.2",
"nyc": "^17.0.0",
"pre-commit": "^1.2.2",
"tape": "^5.0.0",
"typescript": "^5.0.4"
},
"dependencies": {
"reusify": "^1.0.4"
}
}
+346
View File
@@ -0,0 +1,346 @@
'use strict'
/* eslint-disable no-var */
var reusify = require('reusify')
function fastqueue (context, worker, _concurrency) {
if (typeof context === 'function') {
_concurrency = worker
worker = context
context = null
}
if (!(_concurrency >= 1)) {
throw new Error('fastqueue concurrency must be equal to or greater than 1')
}
var cache = reusify(Task)
var queueHead = null
var queueTail = null
var _running = 0
var errorHandler = null
var self = {
push: push,
drain: noop,
saturated: noop,
pause: pause,
paused: false,
get concurrency () {
return _concurrency
},
set concurrency (value) {
if (!(value >= 1)) {
throw new Error('fastqueue concurrency must be equal to or greater than 1')
}
_concurrency = value
if (self.paused) return
for (; queueHead && _running < _concurrency;) {
_running++
release()
}
},
running: running,
resume: resume,
idle: idle,
length: length,
getQueue: getQueue,
unshift: unshift,
empty: noop,
kill: kill,
killAndDrain: killAndDrain,
error: error,
abort: abort
}
return self
function running () {
return _running
}
function pause () {
self.paused = true
}
function length () {
var current = queueHead
var counter = 0
while (current) {
current = current.next
counter++
}
return counter
}
function getQueue () {
var current = queueHead
var tasks = []
while (current) {
tasks.push(current.value)
current = current.next
}
return tasks
}
function resume () {
if (!self.paused) return
self.paused = false
if (queueHead === null) {
_running++
release()
return
}
for (; queueHead && _running < _concurrency;) {
_running++
release()
}
}
function idle () {
return _running === 0 && self.length() === 0
}
function push (value, done) {
var current = cache.get()
current.context = context
current.release = release
current.value = value
current.callback = done || noop
current.errorHandler = errorHandler
if (_running >= _concurrency || self.paused) {
if (queueTail) {
queueTail.next = current
queueTail = current
} else {
queueHead = current
queueTail = current
self.saturated()
}
} else {
_running++
worker.call(context, current.value, current.worked)
}
}
function unshift (value, done) {
var current = cache.get()
current.context = context
current.release = release
current.value = value
current.callback = done || noop
current.errorHandler = errorHandler
if (_running >= _concurrency || self.paused) {
if (queueHead) {
current.next = queueHead
queueHead = current
} else {
queueHead = current
queueTail = current
self.saturated()
}
} else {
_running++
worker.call(context, current.value, current.worked)
}
}
function release (holder) {
if (holder) {
cache.release(holder)
}
var next = queueHead
if (next && _running <= _concurrency) {
if (!self.paused) {
if (queueTail === queueHead) {
queueTail = null
}
queueHead = next.next
next.next = null
worker.call(context, next.value, next.worked)
if (queueTail === null) {
self.empty()
}
} else {
_running--
}
} else if (--_running === 0) {
self.drain()
}
}
function kill () {
queueHead = null
queueTail = null
self.drain = noop
}
function killAndDrain () {
queueHead = null
queueTail = null
self.drain()
self.drain = noop
}
function abort () {
var current = queueHead
queueHead = null
queueTail = null
while (current) {
var next = current.next
var callback = current.callback
var errorHandler = current.errorHandler
var val = current.value
var context = current.context
// Reset the task state
current.value = null
current.callback = noop
current.errorHandler = null
// Call error handler if present
if (errorHandler) {
errorHandler(new Error('abort'), val)
}
// Call callback with error
callback.call(context, new Error('abort'))
// Release the task back to the pool
current.release(current)
current = next
}
self.drain = noop
}
function error (handler) {
errorHandler = handler
}
}
function noop () {}
function Task () {
this.value = null
this.callback = noop
this.next = null
this.release = noop
this.context = null
this.errorHandler = null
var self = this
this.worked = function worked (err, result) {
var callback = self.callback
var errorHandler = self.errorHandler
var val = self.value
self.value = null
self.callback = noop
if (self.errorHandler) {
errorHandler(err, val)
}
callback.call(self.context, err, result)
self.release(self)
}
}
function queueAsPromised (context, worker, _concurrency) {
if (typeof context === 'function') {
_concurrency = worker
worker = context
context = null
}
function asyncWrapper (arg, cb) {
worker.call(this, arg)
.then(function (res) {
cb(null, res)
}, cb)
}
var queue = fastqueue(context, asyncWrapper, _concurrency)
var pushCb = queue.push
var unshiftCb = queue.unshift
queue.push = push
queue.unshift = unshift
queue.drained = drained
return queue
function push (value) {
var p = new Promise(function (resolve, reject) {
pushCb(value, function (err, result) {
if (err) {
reject(err)
return
}
resolve(result)
})
})
// Let's fork the promise chain to
// make the error bubble up to the user but
// not lead to a unhandledRejection
p.catch(noop)
return p
}
function unshift (value) {
var p = new Promise(function (resolve, reject) {
unshiftCb(value, function (err, result) {
if (err) {
reject(err)
return
}
resolve(result)
})
})
// Let's fork the promise chain to
// make the error bubble up to the user but
// not lead to a unhandledRejection
p.catch(noop)
return p
}
function drained () {
var p = new Promise(function (resolve) {
process.nextTick(function () {
if (queue.idle()) {
resolve()
} else {
var previousDrain = queue.drain
queue.drain = function () {
if (typeof previousDrain === 'function') previousDrain()
resolve()
queue.drain = previousDrain
}
}
})
})
return p
}
}
module.exports = fastqueue
module.exports.promise = queueAsPromised
+83
View File
@@ -0,0 +1,83 @@
import * as fastq from '../'
import { promise as queueAsPromised } from '../'
// Basic example
const queue = fastq(worker, 1)
queue.push('world', (err, result) => {
if (err) throw err
console.log('the result is', result)
})
queue.push('push without cb')
queue.concurrency
queue.drain()
queue.empty = () => undefined
console.log('the queue tasks are', queue.getQueue())
queue.idle()
queue.kill()
queue.killAndDrain()
queue.length
queue.pause()
queue.resume()
queue.running()
queue.saturated = () => undefined
queue.unshift('world', (err, result) => {
if (err) throw err
console.log('the result is', result)
})
queue.unshift('unshift without cb')
function worker(task: any, cb: fastq.done) {
cb(null, 'hello ' + task)
}
// Generics example
interface GenericsContext {
base: number;
}
const genericsQueue = fastq<GenericsContext, number, string>({ base: 6 }, genericsWorker, 1)
genericsQueue.push(7, (err, done) => {
if (err) throw err
console.log('the result is', done)
})
genericsQueue.unshift(7, (err, done) => {
if (err) throw err
console.log('the result is', done)
})
function genericsWorker(this: GenericsContext, task: number, cb: fastq.done<string>) {
cb(null, 'the meaning of life is ' + (this.base * task))
}
const queue2 = queueAsPromised(asyncWorker, 1)
async function asyncWorker(task: any) {
return 'hello ' + task
}
async function run () {
await queue.push(42)
await queue.unshift(42)
}
run()
+325
View File
@@ -0,0 +1,325 @@
'use strict'
const test = require('tape')
const buildQueue = require('../').promise
const { promisify } = require('util')
const sleep = promisify(setTimeout)
const immediate = promisify(setImmediate)
test('concurrency', function (t) {
t.plan(2)
t.throws(buildQueue.bind(null, worker, 0))
t.doesNotThrow(buildQueue.bind(null, worker, 1))
async function worker (arg) {
return true
}
})
test('worker execution', async function (t) {
const queue = buildQueue(worker, 1)
const result = await queue.push(42)
t.equal(result, true, 'result matches')
async function worker (arg) {
t.equal(arg, 42)
return true
}
})
test('limit', async function (t) {
const queue = buildQueue(worker, 1)
const [res1, res2] = await Promise.all([queue.push(10), queue.push(0)])
t.equal(res1, 10, 'the result matches')
t.equal(res2, 0, 'the result matches')
async function worker (arg) {
await sleep(arg)
return arg
}
})
test('multiple executions', async function (t) {
const queue = buildQueue(worker, 1)
const toExec = [1, 2, 3, 4, 5]
const expected = ['a', 'b', 'c', 'd', 'e']
let count = 0
await Promise.all(toExec.map(async function (task, i) {
const result = await queue.push(task)
t.equal(result, expected[i], 'the result matches')
}))
async function worker (arg) {
t.equal(arg, toExec[count], 'arg matches')
return expected[count++]
}
})
test('drained', async function (t) {
const queue = buildQueue(worker, 2)
const toExec = new Array(10).fill(10)
let count = 0
async function worker (arg) {
await sleep(arg)
count++
}
toExec.forEach(function (i) {
queue.push(i)
})
await queue.drained()
t.equal(count, toExec.length)
toExec.forEach(function (i) {
queue.push(i)
})
await queue.drained()
t.equal(count, toExec.length * 2)
})
test('drained with exception should not throw', async function (t) {
const queue = buildQueue(worker, 2)
const toExec = new Array(10).fill(10)
async function worker () {
throw new Error('foo')
}
toExec.forEach(function (i) {
queue.push(i)
})
await queue.drained()
})
test('drained with drain function', async function (t) {
let drainCalled = false
const queue = buildQueue(worker, 2)
queue.drain = function () {
drainCalled = true
}
const toExec = new Array(10).fill(10)
let count = 0
async function worker (arg) {
await sleep(arg)
count++
}
toExec.forEach(function () {
queue.push()
})
await queue.drained()
t.equal(count, toExec.length)
t.equal(drainCalled, true)
})
test('drained while idle should resolve', async function (t) {
const queue = buildQueue(worker, 2)
async function worker (arg) {
await sleep(arg)
}
await queue.drained()
})
test('drained while idle should not call the drain function', async function (t) {
let drainCalled = false
const queue = buildQueue(worker, 2)
queue.drain = function () {
drainCalled = true
}
async function worker (arg) {
await sleep(arg)
}
await queue.drained()
t.equal(drainCalled, false)
})
test('set this', async function (t) {
t.plan(1)
const that = {}
const queue = buildQueue(that, worker, 1)
await queue.push(42)
async function worker (arg) {
t.equal(this, that, 'this matches')
}
})
test('unshift', async function (t) {
const queue = buildQueue(worker, 1)
const expected = [1, 2, 3, 4]
await Promise.all([
queue.push(1),
queue.push(4),
queue.unshift(3),
queue.unshift(2)
])
t.is(expected.length, 0)
async function worker (arg) {
t.equal(expected.shift(), arg, 'tasks come in order')
}
})
test('push with worker throwing error', async function (t) {
t.plan(5)
const q = buildQueue(async function (task, cb) {
throw new Error('test error')
}, 1)
q.error(function (err, task) {
t.ok(err instanceof Error, 'global error handler should catch the error')
t.match(err.message, /test error/, 'error message should be "test error"')
t.equal(task, 42, 'The task executed should be passed')
})
try {
await q.push(42)
} catch (err) {
t.ok(err instanceof Error, 'push callback should catch the error')
t.match(err.message, /test error/, 'error message should be "test error"')
}
})
test('unshift with worker throwing error', async function (t) {
t.plan(2)
const q = buildQueue(async function (task, cb) {
throw new Error('test error')
}, 1)
try {
await q.unshift(42)
} catch (err) {
t.ok(err instanceof Error, 'push callback should catch the error')
t.match(err.message, /test error/, 'error message should be "test error"')
}
})
test('no unhandledRejection (push)', async function (t) {
function handleRejection () {
t.fail('unhandledRejection')
}
process.once('unhandledRejection', handleRejection)
const q = buildQueue(async function (task, cb) {
throw new Error('test error')
}, 1)
q.push(42)
await immediate()
process.removeListener('unhandledRejection', handleRejection)
})
test('no unhandledRejection (unshift)', async function (t) {
function handleRejection () {
t.fail('unhandledRejection')
}
process.once('unhandledRejection', handleRejection)
const q = buildQueue(async function (task, cb) {
throw new Error('test error')
}, 1)
q.unshift(42)
await immediate()
process.removeListener('unhandledRejection', handleRejection)
})
test('drained should resolve after async tasks complete', async function (t) {
const logs = []
async function processTask () {
await new Promise(resolve => setTimeout(resolve, 0))
logs.push('processed')
}
const queue = buildQueue(processTask, 1)
queue.drain = () => logs.push('called drain')
queue.drained().then(() => logs.push('drained promise resolved'))
await Promise.all([
queue.push(),
queue.push(),
queue.push()
])
t.deepEqual(logs, [
'processed',
'processed',
'processed',
'called drain',
'drained promise resolved'
], 'events happened in correct order')
})
test('drained should handle undefined drain function', async function (t) {
const queue = buildQueue(worker, 1)
async function worker (arg) {
await sleep(10)
return arg
}
queue.drain = undefined
queue.push(1)
await queue.drained()
t.pass('drained resolved successfully with undefined drain')
})
test('abort rejects all pending promises', async function (t) {
const queue = buildQueue(worker, 1)
const promises = []
let rejectedCount = 0
// Pause queue to prevent tasks from starting
queue.pause()
for (let i = 0; i < 10; i++) {
promises.push(queue.push(i))
}
queue.abort()
// All promises should be rejected
for (const promise of promises) {
try {
await promise
t.fail('promise should have been rejected')
} catch (err) {
t.equal(err.message, 'abort', 'error message is abort')
rejectedCount++
}
}
t.equal(rejectedCount, 10, 'all promises were rejected')
t.equal(queue.length(), 0, 'queue is empty')
async function worker (arg) {
await sleep(500)
return arg
}
})
+733
View File
@@ -0,0 +1,733 @@
'use strict'
/* eslint-disable no-var */
var test = require('tape')
var buildQueue = require('../')
test('concurrency', function (t) {
t.plan(6)
t.throws(buildQueue.bind(null, worker, 0))
t.throws(buildQueue.bind(null, worker, NaN))
t.doesNotThrow(buildQueue.bind(null, worker, 1))
var queue = buildQueue(worker, 1)
t.throws(function () {
queue.concurrency = 0
})
t.throws(function () {
queue.concurrency = NaN
})
t.doesNotThrow(function () {
queue.concurrency = 2
})
function worker (arg, cb) {
cb(null, true)
}
})
test('worker execution', function (t) {
t.plan(3)
var queue = buildQueue(worker, 1)
queue.push(42, function (err, result) {
t.error(err, 'no error')
t.equal(result, true, 'result matches')
})
function worker (arg, cb) {
t.equal(arg, 42)
cb(null, true)
}
})
test('limit', function (t) {
t.plan(4)
var expected = [10, 0]
var queue = buildQueue(worker, 1)
queue.push(10, result)
queue.push(0, result)
function result (err, arg) {
t.error(err, 'no error')
t.equal(arg, expected.shift(), 'the result matches')
}
function worker (arg, cb) {
setTimeout(cb, arg, null, arg)
}
})
test('multiple executions', function (t) {
t.plan(15)
var queue = buildQueue(worker, 1)
var toExec = [1, 2, 3, 4, 5]
var count = 0
toExec.forEach(function (task) {
queue.push(task, done)
})
function done (err, result) {
t.error(err, 'no error')
t.equal(result, toExec[count - 1], 'the result matches')
}
function worker (arg, cb) {
t.equal(arg, toExec[count], 'arg matches')
count++
setImmediate(cb, null, arg)
}
})
test('multiple executions, one after another', function (t) {
t.plan(15)
var queue = buildQueue(worker, 1)
var toExec = [1, 2, 3, 4, 5]
var count = 0
queue.push(toExec[0], done)
function done (err, result) {
t.error(err, 'no error')
t.equal(result, toExec[count - 1], 'the result matches')
if (count < toExec.length) {
queue.push(toExec[count], done)
}
}
function worker (arg, cb) {
t.equal(arg, toExec[count], 'arg matches')
count++
setImmediate(cb, null, arg)
}
})
test('set this', function (t) {
t.plan(3)
var that = {}
var queue = buildQueue(that, worker, 1)
queue.push(42, function (err, result) {
t.error(err, 'no error')
t.equal(this, that, 'this matches')
})
function worker (arg, cb) {
t.equal(this, that, 'this matches')
cb(null, true)
}
})
test('drain', function (t) {
t.plan(4)
var queue = buildQueue(worker, 1)
var worked = false
queue.push(42, function (err, result) {
t.error(err, 'no error')
t.equal(result, true, 'result matches')
})
queue.drain = function () {
t.equal(true, worked, 'drained')
}
function worker (arg, cb) {
t.equal(arg, 42)
worked = true
setImmediate(cb, null, true)
}
})
test('pause && resume', function (t) {
t.plan(13)
var queue = buildQueue(worker, 1)
var worked = false
var expected = [42, 24]
t.notOk(queue.paused, 'it should not be paused')
queue.pause()
queue.push(42, function (err, result) {
t.error(err, 'no error')
t.equal(result, true, 'result matches')
})
queue.push(24, function (err, result) {
t.error(err, 'no error')
t.equal(result, true, 'result matches')
})
t.notOk(worked, 'it should be paused')
t.ok(queue.paused, 'it should be paused')
queue.resume()
queue.pause()
queue.resume()
queue.resume() // second resume is a no-op
function worker (arg, cb) {
t.notOk(queue.paused, 'it should not be paused')
t.ok(queue.running() <= queue.concurrency, 'should respect the concurrency')
t.equal(arg, expected.shift())
worked = true
process.nextTick(function () { cb(null, true) })
}
})
test('pause in flight && resume', function (t) {
t.plan(16)
var queue = buildQueue(worker, 1)
var expected = [42, 24, 12]
t.notOk(queue.paused, 'it should not be paused')
queue.push(42, function (err, result) {
t.error(err, 'no error')
t.equal(result, true, 'result matches')
t.ok(queue.paused, 'it should be paused')
process.nextTick(function () {
queue.resume()
queue.pause()
queue.resume()
})
})
queue.push(24, function (err, result) {
t.error(err, 'no error')
t.equal(result, true, 'result matches')
t.notOk(queue.paused, 'it should not be paused')
})
queue.push(12, function (err, result) {
t.error(err, 'no error')
t.equal(result, true, 'result matches')
t.notOk(queue.paused, 'it should not be paused')
})
queue.pause()
function worker (arg, cb) {
t.ok(queue.running() <= queue.concurrency, 'should respect the concurrency')
t.equal(arg, expected.shift())
process.nextTick(function () { cb(null, true) })
}
})
test('altering concurrency', function (t) {
t.plan(24)
var queue = buildQueue(worker, 1)
queue.push(24, workDone)
queue.push(24, workDone)
queue.push(24, workDone)
queue.pause()
queue.concurrency = 3 // concurrency changes are ignored while paused
queue.concurrency = 2
queue.resume()
t.equal(queue.running(), 2, '2 jobs running')
queue.concurrency = 3
t.equal(queue.running(), 3, '3 jobs running')
queue.concurrency = 1
t.equal(queue.running(), 3, '3 jobs running') // running jobs can't be killed
queue.push(24, workDone)
queue.push(24, workDone)
queue.push(24, workDone)
queue.push(24, workDone)
function workDone (err, result) {
t.error(err, 'no error')
t.equal(result, true, 'result matches')
}
function worker (arg, cb) {
t.ok(queue.running() <= queue.concurrency, 'should respect the concurrency')
setImmediate(function () {
cb(null, true)
})
}
})
test('idle()', function (t) {
t.plan(12)
var queue = buildQueue(worker, 1)
t.ok(queue.idle(), 'queue is idle')
queue.push(42, function (err, result) {
t.error(err, 'no error')
t.equal(result, true, 'result matches')
t.notOk(queue.idle(), 'queue is not idle')
})
queue.push(42, function (err, result) {
t.error(err, 'no error')
t.equal(result, true, 'result matches')
// it will go idle after executing this function
setImmediate(function () {
t.ok(queue.idle(), 'queue is now idle')
})
})
t.notOk(queue.idle(), 'queue is not idle')
function worker (arg, cb) {
t.notOk(queue.idle(), 'queue is not idle')
t.equal(arg, 42)
setImmediate(cb, null, true)
}
})
test('saturated', function (t) {
t.plan(9)
var queue = buildQueue(worker, 1)
var preworked = 0
var worked = 0
queue.saturated = function () {
t.pass('saturated')
t.equal(preworked, 1, 'started 1 task')
t.equal(worked, 0, 'worked zero task')
}
queue.push(42, done)
queue.push(42, done)
function done (err, result) {
t.error(err, 'no error')
t.equal(result, true, 'result matches')
}
function worker (arg, cb) {
t.equal(arg, 42)
preworked++
setImmediate(function () {
worked++
cb(null, true)
})
}
})
test('length', function (t) {
t.plan(7)
var queue = buildQueue(worker, 1)
t.equal(queue.length(), 0, 'nothing waiting')
queue.push(42, done)
t.equal(queue.length(), 0, 'nothing waiting')
queue.push(42, done)
t.equal(queue.length(), 1, 'one task waiting')
queue.push(42, done)
t.equal(queue.length(), 2, 'two tasks waiting')
function done (err, result) {
t.error(err, 'no error')
}
function worker (arg, cb) {
setImmediate(function () {
cb(null, true)
})
}
})
test('getQueue', function (t) {
t.plan(10)
var queue = buildQueue(worker, 1)
t.equal(queue.getQueue().length, 0, 'nothing waiting')
queue.push(42, done)
t.equal(queue.getQueue().length, 0, 'nothing waiting')
queue.push(42, done)
t.equal(queue.getQueue().length, 1, 'one task waiting')
t.equal(queue.getQueue()[0], 42, 'should be equal')
queue.push(43, done)
t.equal(queue.getQueue().length, 2, 'two tasks waiting')
t.equal(queue.getQueue()[0], 42, 'should be equal')
t.equal(queue.getQueue()[1], 43, 'should be equal')
function done (err, result) {
t.error(err, 'no error')
}
function worker (arg, cb) {
setImmediate(function () {
cb(null, true)
})
}
})
test('unshift', function (t) {
t.plan(8)
var queue = buildQueue(worker, 1)
var expected = [1, 2, 3, 4]
queue.push(1, done)
queue.push(4, done)
queue.unshift(3, done)
queue.unshift(2, done)
function done (err, result) {
t.error(err, 'no error')
}
function worker (arg, cb) {
t.equal(expected.shift(), arg, 'tasks come in order')
setImmediate(function () {
cb(null, true)
})
}
})
test('unshift && empty', function (t) {
t.plan(2)
var queue = buildQueue(worker, 1)
var completed = false
queue.pause()
queue.empty = function () {
t.notOk(completed, 'the task has not completed yet')
}
queue.unshift(1, done)
queue.resume()
function done (err, result) {
completed = true
t.error(err, 'no error')
}
function worker (arg, cb) {
setImmediate(function () {
cb(null, true)
})
}
})
test('push && empty', function (t) {
t.plan(2)
var queue = buildQueue(worker, 1)
var completed = false
queue.pause()
queue.empty = function () {
t.notOk(completed, 'the task has not completed yet')
}
queue.push(1, done)
queue.resume()
function done (err, result) {
completed = true
t.error(err, 'no error')
}
function worker (arg, cb) {
setImmediate(function () {
cb(null, true)
})
}
})
test('kill', function (t) {
t.plan(5)
var queue = buildQueue(worker, 1)
var expected = [1]
var predrain = queue.drain
queue.drain = function drain () {
t.fail('drain should never be called')
}
queue.push(1, done)
queue.push(4, done)
queue.unshift(3, done)
queue.unshift(2, done)
queue.kill()
function done (err, result) {
t.error(err, 'no error')
setImmediate(function () {
t.equal(queue.length(), 0, 'no queued tasks')
t.equal(queue.running(), 0, 'no running tasks')
t.equal(queue.drain, predrain, 'drain is back to default')
})
}
function worker (arg, cb) {
t.equal(expected.shift(), arg, 'tasks come in order')
setImmediate(function () {
cb(null, true)
})
}
})
test('killAndDrain', function (t) {
t.plan(6)
var queue = buildQueue(worker, 1)
var expected = [1]
var predrain = queue.drain
queue.drain = function drain () {
t.pass('drain has been called')
}
queue.push(1, done)
queue.push(4, done)
queue.unshift(3, done)
queue.unshift(2, done)
queue.killAndDrain()
function done (err, result) {
t.error(err, 'no error')
setImmediate(function () {
t.equal(queue.length(), 0, 'no queued tasks')
t.equal(queue.running(), 0, 'no running tasks')
t.equal(queue.drain, predrain, 'drain is back to default')
})
}
function worker (arg, cb) {
t.equal(expected.shift(), arg, 'tasks come in order')
setImmediate(function () {
cb(null, true)
})
}
})
test('pause && idle', function (t) {
t.plan(11)
var queue = buildQueue(worker, 1)
var worked = false
t.notOk(queue.paused, 'it should not be paused')
t.ok(queue.idle(), 'should be idle')
queue.pause()
queue.push(42, function (err, result) {
t.error(err, 'no error')
t.equal(result, true, 'result matches')
})
t.notOk(worked, 'it should be paused')
t.ok(queue.paused, 'it should be paused')
t.notOk(queue.idle(), 'should not be idle')
queue.resume()
t.notOk(queue.paused, 'it should not be paused')
t.notOk(queue.idle(), 'it should not be idle')
function worker (arg, cb) {
t.equal(arg, 42)
worked = true
process.nextTick(cb.bind(null, null, true))
process.nextTick(function () {
t.ok(queue.idle(), 'is should be idle')
})
}
})
test('push without cb', function (t) {
t.plan(1)
var queue = buildQueue(worker, 1)
queue.push(42)
function worker (arg, cb) {
t.equal(arg, 42)
cb()
}
})
test('unshift without cb', function (t) {
t.plan(1)
var queue = buildQueue(worker, 1)
queue.unshift(42)
function worker (arg, cb) {
t.equal(arg, 42)
cb()
}
})
test('push with worker throwing error', function (t) {
t.plan(5)
var q = buildQueue(function (task, cb) {
cb(new Error('test error'), null)
}, 1)
q.error(function (err, task) {
t.ok(err instanceof Error, 'global error handler should catch the error')
t.match(err.message, /test error/, 'error message should be "test error"')
t.equal(task, 42, 'The task executed should be passed')
})
q.push(42, function (err) {
t.ok(err instanceof Error, 'push callback should catch the error')
t.match(err.message, /test error/, 'error message should be "test error"')
})
})
test('unshift with worker throwing error', function (t) {
t.plan(5)
var q = buildQueue(function (task, cb) {
cb(new Error('test error'), null)
}, 1)
q.error(function (err, task) {
t.ok(err instanceof Error, 'global error handler should catch the error')
t.match(err.message, /test error/, 'error message should be "test error"')
t.equal(task, 42, 'The task executed should be passed')
})
q.unshift(42, function (err) {
t.ok(err instanceof Error, 'unshift callback should catch the error')
t.match(err.message, /test error/, 'error message should be "test error"')
})
})
test('pause/resume should trigger drain event', function (t) {
t.plan(1)
var queue = buildQueue(worker, 1)
queue.pause()
queue.drain = function () {
t.pass('drain should be called')
}
function worker (arg, cb) {
cb(null, true)
}
queue.resume()
})
test('paused flag', function (t) {
t.plan(2)
var queue = buildQueue(function (arg, cb) {
cb(null)
}, 1)
t.equal(queue.paused, false)
queue.pause()
t.equal(queue.paused, true)
})
test('abort', function (t) {
t.plan(11)
var queue = buildQueue(worker, 1)
var abortedTasks = 0
var predrain = queue.drain
queue.drain = function drain () {
t.fail('drain should never be called')
}
// Pause queue to prevent tasks from starting
queue.pause()
queue.push(1, doneAborted)
queue.push(4, doneAborted)
queue.unshift(3, doneAborted)
queue.unshift(2, doneAborted)
// Abort all queued tasks
queue.abort()
// Verify state after abort
t.equal(queue.length(), 0, 'no queued tasks after abort')
t.equal(queue.drain, predrain, 'drain is back to default')
setImmediate(function () {
t.equal(abortedTasks, 4, 'all queued tasks were aborted')
})
function doneAborted (err) {
t.ok(err, 'error is present')
t.equal(err.message, 'abort', 'error message is abort')
abortedTasks++
}
function worker (arg, cb) {
t.fail('worker should not be called')
setImmediate(function () {
cb(null, true)
})
}
})
test('abort with error handler', function (t) {
t.plan(7)
var queue = buildQueue(worker, 1)
var errorHandlerCalled = 0
queue.error(function (err, task) {
t.equal(err.message, 'abort', 'error handler receives abort error')
t.ok(task !== null, 'error handler receives task value')
errorHandlerCalled++
})
// Pause queue to prevent tasks from starting
queue.pause()
queue.push(1, doneAborted)
queue.push(2, doneAborted)
// Abort all queued tasks
queue.abort()
setImmediate(function () {
t.equal(errorHandlerCalled, 2, 'error handler called for all aborted tasks')
})
function doneAborted (err) {
t.ok(err, 'callback receives error')
}
function worker (arg, cb) {
t.fail('worker should not be called')
setImmediate(function () {
cb(null, true)
})
}
})
+11
View File
@@ -0,0 +1,11 @@
{
"compilerOptions": {
"target": "es6",
"module": "commonjs",
"noEmit": true,
"strict": true
},
"files": [
"./example.ts"
]
}