collabcentral/lib/jobs.js
Joseph B. McQueen 8c3a32f47d Phase R2: extract lib/webex, lib/translation, lib/jobs from index.js
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>
2026-07-02 08:33:42 -04:00

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,
};
}