Server : nginx/1.20.1 System : Linux ccpf-production-2021 5.4.0-148-generic #165-Ubuntu SMP Tue Apr 18 08:53:12 UTC 2023 x86_64 User : forge ( 1000) PHP Version : 7.4.21 Disable Function : pcntl_alarm,pcntl_fork,pcntl_waitpid,pcntl_wait,pcntl_wifexited,pcntl_wifstopped,pcntl_wifsignaled,pcntl_wifcontinued,pcntl_wexitstatus,pcntl_wtermsig,pcntl_wstopsig,pcntl_signal,pcntl_signal_get_handler,pcntl_signal_dispatch,pcntl_get_last_error,pcntl_strerror,pcntl_sigprocmask,pcntl_sigwaitinfo,pcntl_sigtimedwait,pcntl_exec,pcntl_getpriority,pcntl_setpriority,pcntl_async_signals,pcntl_unshare, Directory : /usr/lib/node_modules/pm2/node_modules/culvert/ |
"use strict";
module.exports = makeChannel;
function makeChannel(bufferSize, monitor) {
bufferSize = bufferSize|0;
var dataQueue = [];
var readQueue = [];
var drainList = [];
if (typeof monitor === "string") {
monitor = log(monitor);
}
return {
drain: drain,
put: put,
take: take,
};
function drain(callback) {
if (typeof callback !== "function") {
throw new TypeError("callback must be function");
}
if (dataQueue.length <= bufferSize) return callback();
drainList.push(callback);
}
// Returns true when it's safe to continue without draining
function put(item) {
if (monitor) monitor("put", item);
if (readQueue.length) {
if (monitor) monitor("take", item);
readQueue.shift()(null, item);
}
else {
dataQueue.push(item);
}
return dataQueue.length <= bufferSize;
}
function take(callback) {
if (typeof callback !== "function") {
throw new TypeError("callback must be function");
}
if (dataQueue.length) {
var item = dataQueue.shift();
if (monitor) monitor("take", item);
callback(null, item);
if (dataQueue.length <= bufferSize && drainList.length) {
var list = drainList;
drainList = [];
for (var i = 0; i < list.length; i++) {
list[i]();
}
}
return;
}
readQueue.push(callback);
}
}
function log(name) {
return function (type, value) {
console.info(name, type, value);
};
}