forked from tsai/budibase
9 changed files with 86 additions and 55 deletions
@ -0,0 +1,49 @@ |
|||||
|
const workerFarm = require("worker-farm") |
||||
|
|
||||
|
const ThreadType = { |
||||
|
QUERY: "query", |
||||
|
AUTOMATION: "automation", |
||||
|
} |
||||
|
|
||||
|
function typeToFile(type) { |
||||
|
let filename = null |
||||
|
switch (type) { |
||||
|
case ThreadType.QUERY: |
||||
|
filename = "./query" |
||||
|
break |
||||
|
case ThreadType.AUTOMATION: |
||||
|
filename = "./automation" |
||||
|
break |
||||
|
default: |
||||
|
throw "Unknown thread type" |
||||
|
} |
||||
|
return require.resolve(filename) |
||||
|
} |
||||
|
|
||||
|
class Thread { |
||||
|
constructor(type, opts = { timeoutMs: null, count: 1 }) { |
||||
|
const workerOpts = { |
||||
|
autoStart: true, |
||||
|
maxConcurrentWorkers: opts.count ? opts.count : 1, |
||||
|
} |
||||
|
if (opts.timeoutMs) { |
||||
|
workerOpts.maxCallTime = opts.timeoutMs |
||||
|
} |
||||
|
this.workers = workerFarm(workerOpts, typeToFile(type)) |
||||
|
} |
||||
|
|
||||
|
run(data) { |
||||
|
return new Promise((resolve, reject) => { |
||||
|
this.workers(data, (err, response) => { |
||||
|
if (err) { |
||||
|
reject(err) |
||||
|
} else { |
||||
|
resolve(response) |
||||
|
} |
||||
|
}) |
||||
|
}) |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
module.exports.Thread = Thread |
||||
|
module.exports.ThreadType = ThreadType |
||||
@ -1,5 +1,5 @@ |
|||||
const ScriptRunner = require("../scriptRunner") |
const ScriptRunner = require("../utilities/scriptRunner") |
||||
const { integrations } = require("../../integrations") |
const { integrations } = require("../integrations") |
||||
|
|
||||
function formatResponse(resp) { |
function formatResponse(resp) { |
||||
if (typeof resp === "string") { |
if (typeof resp === "string") { |
||||
@ -1,31 +0,0 @@ |
|||||
const workerFarm = require("worker-farm") |
|
||||
const MAX_WORKER_TIME_MS = 10000 |
|
||||
const workers = workerFarm( |
|
||||
{ |
|
||||
autoStart: true, |
|
||||
maxConcurrentWorkers: 1, |
|
||||
maxCallTime: MAX_WORKER_TIME_MS, |
|
||||
}, |
|
||||
require.resolve("./runner") |
|
||||
) |
|
||||
|
|
||||
function runService(data) { |
|
||||
return new Promise((resolve, reject) => { |
|
||||
workers(data, (err, response) => { |
|
||||
if (err) { |
|
||||
reject(err) |
|
||||
} else { |
|
||||
resolve(response) |
|
||||
} |
|
||||
}) |
|
||||
}) |
|
||||
} |
|
||||
|
|
||||
module.exports = async (datasource, queryVerb, query, transformer) => { |
|
||||
return runService({ |
|
||||
datasource, |
|
||||
queryVerb, |
|
||||
query, |
|
||||
transformer, |
|
||||
}) |
|
||||
} |
|
||||
Loading…
Reference in new issue