mirror of https://github.com/Budibase/budibase.git
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 { integrations } = require("../../integrations") |
|||
const ScriptRunner = require("../utilities/scriptRunner") |
|||
const { integrations } = require("../integrations") |
|||
|
|||
function formatResponse(resp) { |
|||
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