index.js had grown to 1,642 lines / ~50 functions. This peels the three biggest self-contained concerns out into their own modules and wires them back in through a dependency-bag factory so each module stays free of module-level mutable state. lib/webex.js — every webexapis.com round-trip (fetchWithRateLimit, whoAmI, findWebexGroup, getGroupMembers, getPersonInfo, findPersonByEmail, sendDirectMessage, sendMessageWithRetry, sendDirectCard, deleteMessage, refreshToken). Factory closes over token getters and the logger. lib/translation.js — Google Translate fan-out (buildTranslations, translateMessage). Reuses the webex 429 helper so the app has one retry policy. lib/jobs.js — cron-driven pipeline (checkScheduledJobs, processRunningQueue, sendQueueJobMessages, buildJobCompletedCard) plus the two recipient-resolution helpers (collectGroupMembers, buildPeopleList). Factory takes jobs/queue/userPrefs/webex/etc so mutable state stays owned by index.js. index.js: instantiates webex/translator/jobsPipeline once, rewires every call site to go through them, and drops ~715 lines of moved code plus a dead msToTime wrapper. Down from 1,642 → 969 lines, 26 top-level functions instead of 50+. Tests: 53/53 helpers still green. Smoke tested /info, /admin/users, /user/groups/list, /user/groups/find, /jobs/list/all, and the unauth 401 path — all match pre-R2 behavior. Co-authored-by: Cursor <cursoragent@cursor.com>
324 lines
15 KiB
JavaScript
324 lines
15 KiB
JavaScript
// Job queue pipeline: promotes scheduled jobs to running when their time
|
|
// comes, drains the send queue for each running job, records stats, and
|
|
// sends the completion card back to the submitter. Also includes the two
|
|
// helpers that turn form input (a CSV file + a comma-separated group id
|
|
// list) into a deduped recipient list used when a "building" job graduates
|
|
// to a real send.
|
|
//
|
|
// index.js owns the mutable state (jobs object, p-queue instance,
|
|
// userPrefs); this module owns the sequencing logic.
|
|
//
|
|
// createJobsPipeline({ jobs, queue, userPrefs, webex, getBotConfig,
|
|
// getBotProfile, saveJobs, msToTime, logger, fs })
|
|
// returns { checkScheduledJobs, processRunningQueue, buildJobCompletedCard,
|
|
// collectGroupMembers, buildPeopleList }.
|
|
|
|
import defaultFs from 'fs';
|
|
|
|
export function createJobsPipeline(deps) {
|
|
var jobs = deps.jobs;
|
|
var queue = deps.queue;
|
|
var userPrefs = deps.userPrefs || {};
|
|
var webex = deps.webex;
|
|
var getBotConfig = deps.getBotConfig || function () { return null; };
|
|
var getBotProfile = deps.getBotProfile || function () { return null; };
|
|
var saveJobs = deps.saveJobs || function () {};
|
|
var msToTime = deps.msToTime || function (n) { return String(n) + 'ms'; };
|
|
var logger = deps.logger || function () {};
|
|
var fs = deps.fs || defaultFs;
|
|
|
|
// Fires once a minute from the cron. Promotes any scheduled job whose
|
|
// scheduledFor is in the past onto jobs.running, then persists and
|
|
// kicks the running-queue drain. No-ops when there's nothing scheduled.
|
|
function checkScheduledJobs() {
|
|
return new Promise(function (resolve) {
|
|
if (jobs.scheduled.length === 0) return resolve();
|
|
|
|
for (var i = 0; i < jobs.scheduled.length; i++) {
|
|
if (new Date(jobs.scheduled[i].scheduledFor) < new Date()) {
|
|
jobs.running.push(jobs.scheduled[i]);
|
|
delete jobs.scheduled[i];
|
|
}
|
|
}
|
|
jobs.scheduled = jobs.scheduled.filter(function (el) {
|
|
return el != null && el !== undefined && el !== '';
|
|
});
|
|
saveJobs();
|
|
processRunningQueue();
|
|
resolve();
|
|
});
|
|
}
|
|
|
|
// Walks every job currently in jobs.running, assigns them jobIds +
|
|
// start times, and fires them all through the send queue in parallel.
|
|
// Waits for all to settle, cleans up empties, and persists.
|
|
function processRunningQueue() {
|
|
return new Promise(function (resolve, reject) {
|
|
var runningQueuePromises = [];
|
|
for (var i = 0; i < jobs.running.length; i++) {
|
|
jobs.lastJobNumber++;
|
|
jobs.running[i].jobId = jobs.lastJobNumber;
|
|
jobs.running[i].startTime = new Date();
|
|
jobs.running[i].queueNumber = i;
|
|
runningQueuePromises.push(sendQueueJobMessages(i));
|
|
}
|
|
|
|
Promise.allSettled(runningQueuePromises)
|
|
.then(function () {
|
|
logger('processRunningQueue', 'Completed');
|
|
jobs.running = jobs.running.filter(function (el) {
|
|
return el != null && el !== undefined && el !== '';
|
|
});
|
|
saveJobs();
|
|
resolve();
|
|
})
|
|
.catch(reject);
|
|
});
|
|
}
|
|
|
|
// Drains one running job's recipient list through the p-queue. Applies
|
|
// per-recipient language preferences from userPrefs when set, else
|
|
// falls back to english markdown. Records per-recipient timing so the
|
|
// completion card can show hit-rate + average latency. Finally moves
|
|
// the job to jobs.completed and DMs the submitter an Adaptive Card
|
|
// summarizing the run.
|
|
function sendQueueJobMessages(queueNumber) {
|
|
return new Promise(async function (resolve, reject) {
|
|
logger('sendQueueJobMessages',
|
|
'Starting JobId:' + jobs.running[queueNumber].jobId + ' Queue#' + queueNumber);
|
|
|
|
for (var x = 0; x < jobs.running[queueNumber].memberList.length; x++) {
|
|
await queue.onSizeLessThan(2);
|
|
(function capture(job, memberIdx) {
|
|
var member = job.memberList[memberIdx];
|
|
var pref = userPrefs[member.id];
|
|
var markdown = (pref && job.message[pref.language] && job.message[pref.language].markdown)
|
|
? job.message[pref.language].markdown
|
|
: job.message.english.markdown;
|
|
queue.add(function () {
|
|
return sendRunningQueueMessage(
|
|
member.id, markdown, job.imageName, job.queueNumber, memberIdx, job.appName
|
|
);
|
|
});
|
|
})(jobs.running[queueNumber], x);
|
|
}
|
|
|
|
await queue.onIdle();
|
|
var endTime = new Date();
|
|
var startTime = jobs.running[queueNumber].startTime;
|
|
var successCnt = 0;
|
|
var msgTime = 0;
|
|
for (var z = 0; z < jobs.running[queueNumber].memberList.length; z++) {
|
|
if (jobs.running[queueNumber].memberList[z].msgTime) {
|
|
successCnt++;
|
|
msgTime += jobs.running[queueNumber].memberList[z].msgTime;
|
|
}
|
|
}
|
|
jobs.running[queueNumber].endTime = new Date();
|
|
jobs.running[queueNumber].stats = {
|
|
totalTime: msToTime(endTime - startTime),
|
|
webexTime: msToTime(msgTime),
|
|
averageTime: msToTime(msgTime / successCnt),
|
|
succeededMsgs: successCnt,
|
|
totalMsgs: jobs.running[queueNumber].memberList.length,
|
|
succeededMsgPct: (successCnt / jobs.running[queueNumber].memberList.length * 100)
|
|
};
|
|
|
|
var cardDetails = buildJobCompletedCard(jobs.running[queueNumber]);
|
|
jobs.completed.push(jobs.running[queueNumber]);
|
|
logger('sendQueueJobMessages',
|
|
'Completed JobId:' + jobs.running[queueNumber].jobId + ' Queue#' + queueNumber);
|
|
|
|
webex.sendDirectCard(jobs.running[queueNumber].senderId, cardDetails, jobs.running[queueNumber].appName)
|
|
.then(function () {
|
|
delete jobs.running[queueNumber];
|
|
resolve(queueNumber);
|
|
})
|
|
.catch(reject);
|
|
});
|
|
}
|
|
|
|
// Sends one message to one recipient with 429 retry. Records the
|
|
// Webex message result and the wall-clock elapsed time on the member
|
|
// entry so the completion card + job detail view have per-recipient
|
|
// status.
|
|
function sendRunningQueueMessage(toPersonId, message, imageName, jobNumber, memberNumber, appName) {
|
|
return new Promise(async function (resolve, reject) {
|
|
logger('sendRunningQueueMessage', 'Sending to ' + toPersonId);
|
|
var startTime = new Date();
|
|
try {
|
|
const response = await webex.sendMessageWithRetry(toPersonId, message, imageName, appName);
|
|
var body = response.ok ? await response.json() : JSON.parse(await response.text());
|
|
jobs.running[jobNumber].memberList[memberNumber].results = body;
|
|
jobs.running[jobNumber].memberList[memberNumber].msgTime = new Date() - startTime;
|
|
if (response.ok) {
|
|
logger('sendRunningQueueMessage', 'Sent successfully to ' + toPersonId);
|
|
}
|
|
resolve(body);
|
|
} catch (error) {
|
|
logger('sendRunningQueueMessage', 'Failed sending to ' + toPersonId);
|
|
reject(error);
|
|
}
|
|
});
|
|
}
|
|
|
|
// Renders the Adaptive Card + markdown fallback that gets DMed to the
|
|
// submitter when their job finishes. Pulls bot label/avatar from
|
|
// config+profile so it looks branded.
|
|
function buildJobCompletedCard(job) {
|
|
var appName = job.appName;
|
|
var botCfg = getBotConfig(appName) || {};
|
|
var profile = getBotProfile(appName) || {};
|
|
var botLabel = botCfg.label || profile.displayName || appName;
|
|
var botDisplayName = profile.displayName || botLabel;
|
|
var avatarUrl = profile.avatar || null;
|
|
var title = botLabel + ' — Job Summary';
|
|
|
|
var stats = job.stats || {};
|
|
var successPct = (typeof stats.succeededMsgPct === 'number')
|
|
? stats.succeededMsgPct.toFixed(1) + '%'
|
|
: (stats.succeededMsgPct + '%');
|
|
var startedStr = job.startTime ? new Date(job.startTime).toLocaleString() : '—';
|
|
var completedStr = job.endTime ? new Date(job.endTime).toLocaleString() : '—';
|
|
|
|
var headerColumns = [];
|
|
if (avatarUrl) {
|
|
headerColumns.push({
|
|
type: 'Column',
|
|
width: 'auto',
|
|
items: [
|
|
{ type: 'Image', url: avatarUrl, altText: botDisplayName, style: 'Person', size: 'Small', spacing: 'None' }
|
|
]
|
|
});
|
|
}
|
|
headerColumns.push({
|
|
type: 'Column',
|
|
width: 'auto',
|
|
items: [
|
|
{ type: 'TextBlock', text: title, wrap: true, size: 'ExtraLarge', weight: 'Bolder', spacing: 'None', horizontalAlignment: 'Left' }
|
|
],
|
|
verticalContentAlignment: 'Center'
|
|
});
|
|
|
|
var card = {
|
|
type: 'AdaptiveCard',
|
|
'$schema': 'http://adaptivecards.io/schemas/adaptive-card.json',
|
|
version: '1.3',
|
|
body: [
|
|
{ type: 'ColumnSet', columns: headerColumns, spacing: 'None' },
|
|
{ type: 'TextBlock', text: (job.message && job.message.english && job.message.english.markdown) || '', wrap: true },
|
|
{
|
|
type: 'FactSet',
|
|
facts: [
|
|
{ title: 'JobID:', value: (job.jobId !== undefined && job.jobId !== null) ? job.jobId.toString() : '—' },
|
|
{ title: 'Started:', value: startedStr },
|
|
{ title: 'Completed:', value: completedStr },
|
|
{ title: 'Duration:', value: stats.totalTime || '—' },
|
|
{ title: 'Avg Message:', value: stats.averageTime || '—' },
|
|
{ title: 'Webex Time:', value: stats.webexTime || '—' },
|
|
{ title: 'Success', value: (stats.succeededMsgs !== undefined ? (stats.succeededMsgs + '/' + stats.totalMsgs + ' (' + successPct + ')') : '—') }
|
|
],
|
|
separator: true
|
|
}
|
|
]
|
|
};
|
|
|
|
var message = '### ' + title + '\n';
|
|
message += '**JobId:** ' + (job.jobId !== undefined ? job.jobId : '—') + '\n';
|
|
message += '**Started:** ' + startedStr + '\n';
|
|
message += '**Completed:** ' + completedStr + '\n';
|
|
message += '**Duration:** ' + (stats.totalTime || '—') + '\n';
|
|
message += '**Avg Message:** ' + (stats.averageTime || '—') + '\n';
|
|
message += '**Webex Time:** ' + (stats.webexTime || '—') + '\n';
|
|
message += '**Success:** ' + (stats.succeededMsgs !== undefined ? (stats.succeededMsgs + '/' + stats.totalMsgs + ' (' + successPct + ')') : '—') + '\n';
|
|
|
|
return { card: card, message: message };
|
|
}
|
|
|
|
// Walks a comma-separated list of Webex group ids and returns the
|
|
// deduped union of every person in every group, expanding sub-groups
|
|
// transitively. Errors on individual groups are logged and skipped so
|
|
// one bad group id doesn't blow up the whole send. Used by buildJob
|
|
// when the compose form includes group selections.
|
|
function collectGroupMembers(groupList) {
|
|
return new Promise(async function (resolve) {
|
|
var groups = groupList.split(',');
|
|
var memberList = [];
|
|
|
|
for (var group of groups) {
|
|
try {
|
|
var memberCnt = 0;
|
|
var groupCnt = 0;
|
|
var groupData = await webex.getGroupMembers(group);
|
|
|
|
for (var member of groupData.members) {
|
|
memberList.push(member);
|
|
memberCnt++;
|
|
}
|
|
if (groupData.groups.length > 0) {
|
|
for (var subGroup of groupData.groups) {
|
|
logger('collectGroupMembers', 'Found new group: ' + subGroup.name);
|
|
groupCnt++;
|
|
groups.push(subGroup.id);
|
|
}
|
|
}
|
|
logger('collectGroupMembers',
|
|
'Found ' + groupData.displayName + ' Members: ' + memberCnt);
|
|
} catch (error) {
|
|
logger('collectGroupMembers',
|
|
'group ' + group + ' failed: ' + (error && error.message || error));
|
|
}
|
|
}
|
|
|
|
var uniqueMembers = [...new Map(memberList.map(m => [m.id, m])).values()];
|
|
resolve(uniqueMembers);
|
|
});
|
|
}
|
|
|
|
// Reads ./uploads/<fileName> as a newline-delimited list of emails,
|
|
// looks each up in Webex, and returns the deduped recipient array
|
|
// shape the send queue expects. Rejected lookups (bad emails, network
|
|
// hiccups) are logged and dropped — a partial recipient list is
|
|
// better than a whole send failing.
|
|
function buildPeopleList(fileName) {
|
|
return new Promise(function (resolve, reject) {
|
|
if (!fileName) return resolve([]);
|
|
var peoplePromises = [];
|
|
var peopleList = [];
|
|
var csvPeople;
|
|
try {
|
|
csvPeople = fs.readFileSync('./uploads/' + fileName, 'utf-8');
|
|
} catch (err) {
|
|
return reject(err);
|
|
}
|
|
csvPeople.split(/\r?\n/).forEach(function (line) {
|
|
if (line) peoplePromises.push(webex.getPersonInfo(line));
|
|
});
|
|
|
|
Promise.allSettled(peoplePromises).then(function (peopleData) {
|
|
for (var person of peopleData) {
|
|
if (person.status === 'fulfilled' && person.value) {
|
|
logger('buildPeopleList', 'Found ' + person.value.displayName);
|
|
peopleList.push({
|
|
id: person.value.id,
|
|
type: 'peopleList',
|
|
displayName: person.value.displayName
|
|
});
|
|
} else if (person.status === 'rejected') {
|
|
logger('buildPeopleList',
|
|
'lookup failed: ' + (person.reason && person.reason.message || person.reason));
|
|
}
|
|
}
|
|
resolve(peopleList);
|
|
});
|
|
});
|
|
}
|
|
|
|
return {
|
|
checkScheduledJobs,
|
|
processRunningQueue,
|
|
buildJobCompletedCard,
|
|
collectGroupMembers,
|
|
buildPeopleList,
|
|
};
|
|
}
|