// Jobs endpoints — the core of the app: // POST /CollabCentral/:app/jobs/:action — edit / runNow / schedule / // cancel // GET /CollabCentral/:app/jobs/list/:scope // GET /CollabCentral/:app/jobs/detail/:jobId // // registerJobsRoutes(app, { jobs, isAuthorized, helpers, webex, // translator, jobsPipeline, saveConfig, upload, // sharp, uuid, fs, logger }) // // Building (draft) jobs are keyed by cookieId + appName so a user // authorized on multiple bots can have one independent draft per bot. export function registerJobsRoutes(app, deps) { var jobs = deps.jobs; var isAuthorized = deps.isAuthorized; var helpers = deps.helpers; var webex = deps.webex; var translator = deps.translator; var jobsPipeline = deps.jobsPipeline; var saveConfig = deps.saveConfig; var upload = deps.upload; var sharp = deps.sharp; var uuid = deps.uuid; var fs = deps.fs; var logger = deps.logger || function () {}; function buildingKey(req) { return helpers.buildingKey(req.cookies && req.cookies.id, req.params.app); } // Takes the multipart/form-data submitted by the compose page and // fans it into a `building` job entry: any uploaded image gets JPEG- // compressed to ./uploads/, any uploaded CSV gets parsed into a // recipient list, and any group ids get expanded through Webex. // Resolves with the fully-formed building job entry once every side // effect settles. function buildJob(req, res) { return new Promise(function (resolve) { res.setTimeout(3000000); var memberList = []; var memberListPromises = []; upload(req, res, async function (err) { jobs.building[buildingKey(req)] = { appName: req.body.appName, senderId: req.cookies.id, senderDisplayName: req.cookies.displayName, groups: req.body.groups, message: { raw: req.body.message, english: {} }, memberList: [] }; if (!err && req.files) { if (req.files['uploadImage']) { for (var file of req.files['uploadImage']) { let fileName = uuid() + '.jpeg'; await sharp(file.buffer).jpeg({ quality: 40 }).toFile('./uploads/' + fileName) .catch(err => logger('buildJob', 'sharp error: ' + err)); jobs.building[buildingKey(req)].imageName = fileName; logger('buildJob', 'uploadImage ' + file.originalname + ' to ' + fileName); } } if (req.files['uploadCSV']) { for (var file of req.files['uploadCSV']) { let fileName = uuid() + '.csv'; fs.writeFileSync('./uploads/' + fileName, file.buffer); jobs.building[buildingKey(req)].csv = fileName; logger('buildJob', 'uploadCSV ' + file.originalname + ' to ' + fileName); memberListPromises.push(jobsPipeline.buildPeopleList(fileName)); } } } if (req.body.groups) { logger('buildJob', 'Locating groups: ' + req.body.groups); memberListPromises.push(jobsPipeline.collectGroupMembers(req.body.groups)); } if (req.body.scheduledFor) { logger('buildJob', 'Scheduling for ' + new Date(req.body.scheduledFor).toLocaleString()); jobs.building[buildingKey(req)].scheduledFor = req.body.scheduledFor; saveConfig(jobs, './config/jobs.json'); } Promise.allSettled(memberListPromises).then(function (peoplePromises) { for (var peopleList of peoplePromises) { if (peopleList.status === 'fulfilled') { for (var person of peopleList.value) { memberList.push(person); } } } var uniqueMembers = [...new Map(memberList.map(m => [m.id, m])).values()]; jobs.building[buildingKey(req)].memberList = uniqueMembers; saveConfig(jobs, './config/jobs.json'); resolve(jobs.building[buildingKey(req)]); }); }); }); } app.post('/CollabCentral/:app/jobs/:action', function (req, res) { if (!isAuthorized(req.params.app, req.cookies.id)) { return res.status(401).send('You are not authorized.'); } req.setTimeout(3000000); if (req.params.action === 'edit') { logger('apiEndpoint(' + req.params.app + ')', 'POST /jobs/edit'); return buildJob(req, res) .then(job => { webex.sendDirectMessage(job.senderId, job.message.raw, job.imageName, req.params.app) .then(async function (result) { jobs.building[buildingKey(req)].message.english = { text: result.text, markdown: result.markdown, html: result.html }; // Remember the preview DM's Webex message id // so /jobs/cancel can retract it and // /jobs/runNow can treat it as "already // delivered" (future). if (result && result.id) { jobs.building[buildingKey(req)].previewMessageId = result.id; } await translator.buildTranslations(jobs.building[buildingKey(req)].message.english) .then(t => { jobs.building[buildingKey(req)].message = t; }) .catch(error => logger('jobs/edit', 'Error translating: ' + error)); res.status(200).send(job); saveConfig(jobs, './config/jobs.json'); }) .catch(error => logger('jobs/edit', 'Error sending message: ' + error)); }) .catch(error => logger('jobs/edit', 'Error buildJob: ' + error)); } if (req.params.action === 'runNow') { logger('apiEndpoint(' + req.params.app + ')', 'POST /jobs/runNow'); jobs.running.push(jobs.building[buildingKey(req)]); delete jobs.building[buildingKey(req)]; jobsPipeline.processRunningQueue() .catch(error => logger('jobs/runNow', 'sendRunningMessages error: ' + error)); return res.status(204).redirect('/CollabCentral/' + req.params.app + '/monitorJobs.html'); } if (req.params.action === 'schedule') { logger('apiEndpoint(' + req.params.app + ')', 'POST /jobs/schedule'); jobs.scheduled.push(jobs.building[buildingKey(req)]); delete jobs.building[buildingKey(req)]; saveConfig(jobs, './config/jobs.json'); return res.status(204).redirect('/CollabCentral/' + req.params.app + '/monitorJobs.html'); } if (req.params.action === 'cancel') { // Abandon the caller's in-progress "building" job for this // bot. A building entry exists as soon as the user clicks // "Review & send" (POST /jobs/edit) and lingers on the // server until the user promotes it to running/scheduled — // or until now, where this route just drops it on the // floor. Idempotent: 200 whether or not a building entry // actually existed for (personId, app), so the client can // always fire this on "cancel" without needing to check // state first. // // If the building job has a previewMessageId (captured // during /jobs/edit), also retract that DM from the // sender's own space so the review message doesn't hang // around after they discarded the draft. Best-effort — a // failed Webex delete is logged but does not fail the // cancel. logger('apiEndpoint(' + req.params.app + ')', 'POST /jobs/cancel'); var key = buildingKey(req); var job = jobs.building[key]; var hadJob = !!job; if (hadJob) { var previewId = job.previewMessageId; delete jobs.building[key]; try { saveConfig(jobs, './config/jobs.json'); } catch (err) { logger('apiEndpoint(' + req.params.app + ')', 'jobs/cancel save failed: ' + (err && err.message || err)); return res.status(500).send('Failed to cancel job.'); } if (previewId) { // Fire-and-forget; we've already committed the // cancel to disk and we don't want a slow Webex // round trip to hold the response. webex.deleteMessage(previewId, req.params.app); } } return res.status(200).send({ cancelled: hadJob }); } return res.status(404).send('Unknown jobs action.'); }); app.get('/CollabCentral/:app/jobs/list/:scope', function (req, res) { if (!req.cookies || !isAuthorized(req.params.app, req.cookies.id)) { return res.status(401).send('You are not authorized.'); } var appName = req.params.app; if (req.params.scope === 'all') { logger('apiEndpoint(' + appName + ')', 'GET /jobs/list/all'); return res.status(200).send({ building: jobs.building[buildingKey(req)], running: helpers.jobsForApp(jobs.running, appName), scheduled: helpers.jobsForApp(jobs.scheduled, appName), completed: helpers.jobsForApp(jobs.completed, appName) }); } if (req.params.scope === 'building') { logger('apiEndpoint(' + appName + ')', 'GET /jobs/list/building'); return res.status(200).send(jobs.building[buildingKey(req)]); } if (req.params.scope === 'running') { logger('apiEndpoint(' + appName + ')', 'GET /jobs/list/running'); var runningJobs = []; for (var job of helpers.jobsForApp(jobs.running, appName)) { var completedSends = 0; for (var member of (job.memberList || [])) { if (member.results) completedSends++; } runningJobs.push({ jobId: job.jobId, senderDisplayName: job.senderDisplayName, appName: job.appName, startTime: job.startTime, totalRecipients: (job.memberList || []).length, completedRecipients: completedSends }); } return res.status(200).send(runningJobs); } if (req.params.scope === 'scheduled') { logger('apiEndpoint(' + appName + ')', 'GET /jobs/list/scheduled'); var scheduledJobs = []; for (var job of helpers.jobsForApp(jobs.scheduled, appName)) { scheduledJobs.push({ senderDisplayName: job.senderDisplayName, appName: job.appName, message: job.message && job.message.english && job.message.english.html, totalRecipients: (job.memberList || []).length, scheduledFor: job.scheduledFor }); } return res.status(200).send(scheduledJobs); } if (req.params.scope === 'completed') { logger('apiEndpoint(' + appName + ')', 'GET /jobs/list/completed'); return res.status(200).send(helpers.jobsForApp(jobs.completed, appName)); } return res.status(404).send(); }); // Returns the full job object for a given jobId, but only if the // job belongs to the requested bot. Looks across // completed/running/scheduled so the same URL keeps working as a // job moves through its lifecycle. app.get('/CollabCentral/:app/jobs/detail/:jobId', function (req, res) { if (!req.cookies || !isAuthorized(req.params.app, req.cookies.id)) { return res.status(401).send('You are not authorized.'); } var appName = req.params.app; var jobId = req.params.jobId; logger('apiEndpoint(' + appName + ')', 'GET /jobs/detail/' + jobId); var all = (jobs.completed || []).concat(jobs.running || [], jobs.scheduled || []); var match = all.find(function (j) { return j && j.appName === appName && String(j.jobId) === String(jobId); }); if (!match) return res.status(404).send('Job not found.'); res.status(200).send(match); }); }