Async queue with configurable concurrency. Callback pattern with onData/onEnd.
Module queue | Source packages/front/fw/src/io/utils/queue.js | Deps none | Worker-safe yes
Resolve
const Queue = runtime.resolve('queue');
// Returns: Queue constructor (class)
API
const q = new Queue(concurrency);
| Param |
Type |
Default |
Description |
concurrency |
number|boolean |
1 |
Max simultaneous jobs. true → 1024 |
Properties
| Property |
Type |
Description |
onData |
Function |
Must be set. (job, end) => void — processes a job, calls end() when done |
onEnd |
Function |
Must be set. (error) => void — called when the queue is closed |
length |
number |
Number of jobs waiting |
running |
number |
Number of jobs in progress |
stopped |
boolean |
true after stop() |
error |
any |
Error passed to stop() or onData end(error) |
Methods
| Method |
Signature |
Description |
push |
(job: any) => void |
Adds a job |
concat |
(jobs: any[]) => void |
Adds multiple jobs |
end |
(error?) => void |
Signals EOF — no more jobs to come |
stop |
(error?) => void |
Stops processing immediately |
Examples
Sequential (concurrency 1)
const Queue = runtime.resolve('queue');
const q = new Queue(1);
q.onData = function(item, done) {
// Async processing
fetch(item.url)
.then(r => r.json())
.then(data => {
console.log(data);
done(); // job done — triggers the next one
})
.catch(err => done(err)); // error → closes the queue
};
q.onEnd = function(error) {
if (error) console.error('Queue error:', error);
else console.log('All done');
};
q.push({ url: '/api/1' });
q.push({ url: '/api/2' });
q.push({ url: '/api/3' });
q.end(); // no more jobs to come
Parallel (concurrency 4)
const q = new Queue(4);
q.onData = (chunk, done) => {
processChunk(chunk).then(() => done());
};
q.onEnd = () => console.log('Batch complete');
q.concat(chunksArray);
q.end();
Emergency stop
q.stop(new Error('Cancelled by user'));
// → onEnd(error) called immediately
Worker Usage
const worker = fw.createWorker(
function ({ libs, args }) {
const Queue = libs.queue;
const q = new Queue(2);
q.onData = (url, done) => {
fetch(url).then(r => r.json())
.then(data => { self.postMessage(data); done(); })
.catch(err => done(err));
};
q.onEnd = (error) => {
if (error) self.postMessage({ error: error.message });
};
q.concat(args.urls);
q.end();
},
{ dependencies: ['queue'], args: { urls: ['/api/1', '/api/2', '/api/3'] } }
);
Notes
onData and onEnd must be set before the first push — otherwise Error at runtime.
done() in onData must be called exactly once — multiple calls → Error.
push() after end() → Error.
concurrency = true → 1024 (near-parallel execution).
- Internal state machine:
PROCESSING | EOF | CLOSING | CLOSED.
See also
- workers — similar pattern for parallel workers