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