diff --git a/index.js b/index.js index 6302519..d68b55a 100644 --- a/index.js +++ b/index.js @@ -13,8 +13,10 @@ import FormData from 'form-data'; import cron from "node-cron"; import PQueue from 'p-queue'; import * as helpers from './lib/helpers.js'; -//import pLimit from "p-limit"; -//const limit = pLimit(10); +import { createWebexClient } from './lib/webex.js'; +import { createTranslator } from './lib/translation.js'; +import { createJobsPipeline } from './lib/jobs.js'; + const queue = new PQueue({ concurrency: 10 }); //Load the config files from storage. @@ -103,7 +105,7 @@ async function loadBotProfile(appName) { var token = getBotToken(appName); if (!token) return; try { - var response = await fetchWithRateLimit('https://webexapis.com/v1/people/me', { + var response = await webex.fetchWithRateLimit('https://webexapis.com/v1/people/me', { method: 'GET', headers: { 'Authorization': 'Bearer ' + token } }); @@ -149,7 +151,7 @@ function refreshGroupsCache() { if (groupsCache.refreshing) return groupsCache.refreshing; var startedAt = Date.now(); logger('groupsCache', 'Refreshing group list from Webex...'); - groupsCache.refreshing = findWebexGroup() + groupsCache.refreshing = webex.findWebexGroup() .then(function (groups) { groupsCache.data = groups; groupsCache.lastRefreshed = new Date(); @@ -246,6 +248,39 @@ function getOAuthRedirectUri(appName) { return helpers.getOAuthRedirectUri(process.env.OAUTH_CALLBACK_URL_TEMPLATE, appName); } +// --------------------------------------------------------------------------- +// Wire the lib/* modules. Each takes a dependency bag so the modules +// themselves stay free of module-level mutable state — everything mutable +// (tokens, jobs, queue) still lives here in index.js. +// --------------------------------------------------------------------------- +var webex = createWebexClient({ + getBotToken: getBotToken, + getServiceAccountAccessToken: getServiceAccountAccessToken, + getServiceAccountRefreshToken: function () { return serviceAccountToken && serviceAccountToken.refresh_token; }, + saveServiceAccountToken: saveServiceAccountToken, + env: process.env, + logger: logger, +}); + +var translator = createTranslator({ + languages: config.languages || [], + apiKey: process.env.GOOGLE_TRANSLATE_API_KEY || '', + fetchWithRateLimit: webex.fetchWithRateLimit, + logger: logger, +}); + +var jobsPipeline = createJobsPipeline({ + jobs: jobs, + queue: queue, + userPrefs: userPrefs, + webex: webex, + getBotConfig: function (appName) { return helpers.getBotConfig(config, botTokens, appName); }, + getBotProfile: getBotProfile, + saveJobs: function () { saveConfig(jobs, './config/jobs.json'); }, + msToTime: helpers.msToTime, + logger: logger, +}); + //Load Express Server var app = express(); app.use(bodyParser.json()); @@ -293,13 +328,13 @@ cron.schedule('0 * * * * *', () => { logger("refreshToken", "Renew by: " + new Date(renewBy).toLocaleString()) logger("refreshToken", "Renew after: " + new Date(new Date(renewBy) - 7200000).toLocaleString()); logger("refreshToken", "Refreshing token.") - refreshToken() + webex.refreshToken() .then(response => { logger("tokenRefresh", response); }) .catch(error => logger("refreshToken", error)); } - checkScheduledJobs() + jobsPipeline.checkScheduledJobs() .catch(error => logger("checkScheduledJobs", error)); }, { timezone: cronTimezone }) @@ -357,7 +392,7 @@ app.post('/CollabCentral/:app/jobs/:action', (req, res) => { logger("apiEndpoint(" + req.params.app + ")", "POST /jobs/edit"); buildJob(req, res) .then(job => { - sendDirectMessage(job.senderId, job.message.raw, job.imageName, req.params.app) + 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, @@ -371,7 +406,7 @@ app.post('/CollabCentral/:app/jobs/:action', (req, res) => { jobs.building[buildingKey(req)].previewMessageId = result.id; } - await buildTranslations(jobs.building[buildingKey(req)].message.english) + await translator.buildTranslations(jobs.building[buildingKey(req)].message.english) .then(result => { jobs.building[buildingKey(req)].message = result; }) @@ -387,7 +422,7 @@ app.post('/CollabCentral/:app/jobs/:action', (req, res) => { logger("apiEndpoint(" + req.params.app + ")", "POST /jobs/runNow"); jobs.running.push(jobs.building[buildingKey(req)]); delete jobs.building[buildingKey(req)]; - processRunningQueue() + jobsPipeline.processRunningQueue() .catch(error => console.log("sendRunningMessages error: " + error)) res.status(204).redirect("/CollabCentral/" + req.params.app + "/monitorJobs.html"); } else if (req.params.action == "schedule") { @@ -429,7 +464,7 @@ app.post('/CollabCentral/:app/jobs/:action', (req, res) => { // 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. - deleteWebexMessage(previewId, req.params.app); + webex.deleteMessage(previewId, req.params.app); } } return res.status(200).send({ cancelled: hadJob }); @@ -675,7 +710,7 @@ app.post('/CollabCentral/:app/admin/users', async function (req, res) { var person; try { - person = await findPersonByEmail(email); + person = await webex.findPersonByEmail(email); } catch (err) { logger('apiEndpoint(' + appName + ')', 'admin/users lookup failed for "' + email + '": ' + (err && err.message || err)); @@ -770,7 +805,7 @@ app.get(`/CollabCentral/:app/oauth`, async function (req, res) { var response; try { - response = await fetchWithRateLimit(url, { method: 'POST', body: params }); + response = await webex.fetchWithRateLimit(url, { method: 'POST', body: params }); } catch (error) { logger("oauth", "Token exchange fetch failed: " + (error && error.message || error)); return res.status(502).send("Failed to reach Webex to complete OAuth."); @@ -785,7 +820,7 @@ app.get(`/CollabCentral/:app/oauth`, async function (req, res) { var jsonData = await response.json(); var whoami; try { - whoami = await whoAmI(jsonData.access_token); + whoami = await webex.whoAmI(jsonData.access_token); } catch (err) { logger("oauth", "whoAmI failed: " + (err && err.message || err)); return res.status(502).send("Failed to read your Webex profile after OAuth."); @@ -844,14 +879,14 @@ function buildJob(req, res) { fs.writeFileSync('./uploads/' + fileName, file.buffer) jobs.building[buildingKey(req)].csv = fileName; logger("buildJob", "uploadCSV " + file.originalname + " to " + fileName) - memberListPromises.push(buildPeopleList(fileName)); + memberListPromises.push(jobsPipeline.buildPeopleList(fileName)); } } } if (req.body.groups) { logger("buildJob", "Locating groups: " + req.body.groups); - memberListPromises.push(collectGroupMembers(req.body.groups)); + memberListPromises.push(jobsPipeline.collectGroupMembers(req.body.groups)); } if (req.body.scheduledFor) { @@ -882,751 +917,6 @@ function buildJob(req, res) { }) } -function collectGroupMembers(groupList) { - return new Promise(async function (resolve, reject) { - - var groups = groupList.split(","); - var memberList = []; - - for (var group of groups) { - try { - var memberCnt = 0; - var groupCnt = 0; - var groupData = await 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) { - console.log("Error groupData: " + error) - } - } - - //console.log("Groups: " + groupList); - var uniqueMembers = [...new Map(memberList.map((m) => [m.id, m])).values()]; - resolve(uniqueMembers); - }) -} - -function getGroupMembers(groupId) { - return new Promise(async function (resolve, reject) { - var groupMembers = { "members": [], "groups": [] }; - var startIndex = 1; - var count = 500; - var memberSize = 10000; - - var requestOptions = { - method: 'GET', - headers: { - 'Authorization': `Bearer ` + getServiceAccountAccessToken(), - 'Content-Type': `application/json` - } - }; - - while (memberSize > (groupMembers.members.length + groupMembers.groups.length)) { - var findGroupByIdUrl = 'https://webexapis.com/v1/groups/' + groupId + '/members?startIndex=' + startIndex + '&count=' + count; - try { - var groupResponse = await fetchWithRateLimit(findGroupByIdUrl, requestOptions); - if (groupResponse.ok) { - var groupData = await groupResponse.json(); - groupMembers.displayName = groupData.displayName; - memberSize = groupData.memberSize; - for (var member of groupData.members) { - if (member.type == "group") { - - groupMembers.groups.push({ "id": member.id, "name": member.displayName }) - } else { - groupMembers.members.push(member) - } - } - startIndex = startIndex + count; - } else { - memberSize = 0; - reject(groupResponse.status + ": " + groupResponse.statusText); - } - } catch (error) { - console.log("Error findGroupById: " + error); - } - } - resolve(groupMembers); - }) -} - -function findWebexGroup() { - return new Promise(async function (resolve, reject) { - - var groups = []; - var startIndex = 1; - var count = 500; - var groupSize = 10000; - - - var url = new URL("https://webexapis.com/v1/groups") - //var params = new URLSearchParams(); - //url.searchParams.append("filter", "displayName contains " + searchTerm); - url.searchParams.append("attributes", "displayName"); - url.searchParams.append("sortBy", "displayName"); - url.searchParams.append("sortOrder", "ascending"); - url.searchParams.append("includeMembers", "false"); - url.searchParams.append("startIndex", startIndex); - url.searchParams.append("count", count); - //console.log(params.toString()) - - var requestOptions = { - method: 'GET', - headers: { - 'Authorization': `Bearer ` + getServiceAccountAccessToken(), - 'Content-Type': `application/json` - } - }; - while (groupSize > groups.length) { - try { - var fetchedGroups = await fetchWithRateLimit(url, requestOptions); - if (fetchedGroups.ok) { - var groupData = await fetchedGroups.json(); - groupSize = groupData.totalResults; - for (var group of groupData.groups) { - groups.push(group); - } - startIndex = startIndex + count; - url.searchParams.delete("startIndex"); - url.searchParams.append("startIndex", startIndex); - } else { - return reject(fetchedGroups.status + ": " + fetchedGroups.statusText); - } - } catch (error) { - logger('findWebexGroup', 'page fetch failed: ' + (error && error.message || error)); - return reject(error); - } - } - resolve(groups); - }) -} -function whoAmI(bearerToken) { - return new Promise(async function (resolve, reject) { - logger("whoAmI", "Person is having an existential crisis of identity.") - var myHeaders = { - "Authorization": "Bearer " + bearerToken, - "Content-Type": "application/json" - } - var requestOptions = { - method: 'GET', - headers: myHeaders, - redirect: 'follow' - }; - - fetchWithRateLimit("https://webexapis.com/v1/people/me", requestOptions) - .then(async function (response) { - var userInfo = await response.json(); - resolve(userInfo); - }) - .catch(error => reject(error)); - }) -} - -function buildPeopleList(fileName) { - return new Promise(async function (resolve, reject) { - var peoplePromises = []; - var peopleList = []; - if (fileName) { - var csvPeople = fs.readFileSync("./uploads/" + fileName, 'utf-8'); - console.info(csvPeople); - csvPeople.split(/\r?\n/).forEach(line => { - peoplePromises.push(getPersonInfo(line)); - - }) - - Promise.allSettled(peoplePromises) - .then(peopleData => { - //console.log(JSON.stringify(peopleData)) - var cnt = 0; - for (var person of peopleData) { - console.log(JSON.stringify(person)) - cnt++; - if (person.status == "fulfilled") { - if (person.value) { - console.log("Found " + person.value.displayName) - peopleList.push({ - "id": person.value.id, - "type": "peopleList", - "displayName": person.value.displayName - }) - } else { - console.log(JSON.stringify(person)) - } - } else { - console.log(cnt) - console.log(JSON.stringify(person)) - } - } - console.log(JSON.stringify(peopleList)) - resolve(peopleList); - }) - .catch(error => reject(error)); - } - - }) -} - -function getPersonInfo(email) { - return new Promise(async function (resolve, reject) { - var myHeaders = { - "Authorization": "Bearer " + getServiceAccountAccessToken(), - "Content-Type": "application/json" - } - var requestOptions = { - method: 'GET', - headers: myHeaders, - redirect: 'follow' - }; - - fetchWithRateLimit("https://webexapis.com/v1/people?email=" + encodeURIComponent(email) + "&max=100", requestOptions) - .then(async function (response) { - console.log(email + ": " + response.status + ":" + response.statusText) - if (response.ok) { - var userInfo = await response.json(); - resolve(userInfo.items[0]); - } else { - console.log(response.status + ": " + response.statusText) - try { - var stuff = await response.text(); - } catch (error) { - console.log(error) - console.log(stuff) - } - reject(); - } - - }) - .catch(error => { - - reject(error) - }); - - - }) -} - -function checkScheduledJobs() { - return new Promise(async function (resolve, reject) { - - if (jobs.scheduled.length > 0) { - 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]; - } else { } - } - - jobs.scheduled = jobs.scheduled.filter(elements => { - return (elements != null && elements !== undefined && elements !== ""); - }); - saveConfig(jobs, "./config/jobs.json"); - processRunningQueue(); - } else { - //console.log("No jobs scheduled."); - } - - }) - -} - -function processRunningQueue() { - return new Promise(async 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(promises => { - - logger("processRunningQueue", "Completed"); - jobs.running = jobs.running.filter(elements => { return (elements != null && elements !== undefined && elements !== ""); }); - saveConfig(jobs, "./config/jobs.json") - resolve(); - }) - .catch(error => reject(error)); - }) -} - -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); - if (userPrefs[jobs.running[queueNumber].memberList[x].id]) { - queue.add(() => sendRunningQueueMessage(jobs.running[queueNumber].memberList[x].id, jobs.running[queueNumber].message[userPrefs[jobs.running[queueNumber].memberList[x].id].language].markdown, jobs.running[queueNumber].imageName, queueNumber, x, jobs.running[queueNumber].appName)); - } else { - queue.add(() => sendRunningQueueMessage(jobs.running[queueNumber].memberList[x].id, jobs.running[queueNumber].message.english.markdown, jobs.running[queueNumber].imageName, queueNumber, x, jobs.running[queueNumber].appName)); - } - - } - - //console.log(membershipsPromise); - 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) - } - - jobs.running[queueNumber].endTime = new Date(); - var cardDetails = buildJobCompletedCard(jobs.running[queueNumber]); - jobs.completed.push(jobs.running[queueNumber]); - logger("sendQueueJobMessages", "Completed JobId:" + jobs.running[queueNumber].jobId + " Queue#" + queueNumber) - - await sendDirectCard(jobs.running[queueNumber].senderId, cardDetails, jobs.running[queueNumber].appName) - .then(msgResults => { - delete jobs.running[queueNumber]; - }) - .catch(error => { - //jobs.running[i].memberList[x].error = error; - reject(error); - }) - - resolve(queueNumber); - - }) -} -function sendRunningQueueMessage(toPersonId, message, imageName, jobNumber, memberNumber, appName) { - return new Promise(async function (resolve, reject) { - //logger("sendRunningQueueMessage", toPersonId + "\n" + message + "\n" + imageName + "\n" + jobNumber + "\n" + memberNumber + "\n" + appName) - logger("sendRunningQueueMessage", "Sending to " + toPersonId); - var startTime = new Date(); - - var myHeaders = new Headers(); - myHeaders.append("Authorization", "Bearer " + getBotToken(appName)); - var formdata = new FormData(); - formdata.append("toPersonId", toPersonId); - formdata.append("markdown", message); - - if (imageName) { - var stream = fs.createReadStream("./uploads/" + imageName); - //console.log(file); - formdata.append("files", stream, imageName); - }; - - //var myBody = JSON.stringify(myBodyJSON); - var requestOptions = { - method: 'POST', - headers: myHeaders, - body: formdata, - redirect: 'follow' - }; - - try { - - const response = await fetchAndRetryIfNecessary(async () => ( - await fetch("https://webexapis.com/v1/messages", requestOptions) - )) - - if (response.ok) { - var result = await response.json(); - jobs.running[jobNumber].memberList[memberNumber].results = result; - jobs.running[jobNumber].memberList[memberNumber].msgTime = new Date() - startTime; - logger("sendRunningQueueMessage", "Sent successfully to " + toPersonId); - resolve(result) - } else { - var result = await response.text(); - jobs.running[jobNumber].memberList[memberNumber].results = JSON.parse(result); - jobs.running[jobNumber].memberList[memberNumber].msgTime = new Date() - startTime; - resolve(JSON.parse(result)); - } - - } catch (error) { - logger("sendRunningQueueMessage", "Failed sending to " + toPersonId); - reject(error) - } - }) -} - -function sleep(milliseconds) { - return new Promise((resolve) => setTimeout(resolve, milliseconds)) -} - -function buildJobCompletedCard(job) { - var appName = job.appName; - var botCfg = (config.webex && config.webex.bot && config.webex.bot[appName]) || {}; - var profile = getBotProfile(appName) || {}; - - // Per-bot identity: prefer the live profile from Webex /people/me, fall - // back to the configured label / appName so the card renders even before - // the profile loader finishes. - 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() : "—"; - - // Header columns: bot avatar (if known) + job summary title. - 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 jobCompletedCard = { - "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": jobCompletedCard, "message": message }; -} - -async function fetchWithRateLimit(url, requestOptions) { - const response = await fetch(url, requestOptions) - - if (response.status === 429) { - //console.log(response.headers) - const secondsToWait = Number(response.headers.get('retry-after')) - console.log("Waiting for " + secondsToWait.toString() + " due to 429 message: " + url); - console.log(JSON.stringify(requestOptions)); - await new Promise(resolve => setTimeout(resolve, secondsToWait * 1000)) - console.log("Finished waiting for " + secondsToWait.toString() + " due to 429 message: " + url); - console.log(JSON.stringify(requestOptions)) - return await fetchWithRateLimit(url, requestOptions) - } - return response -} - -function getMillisToSleep(retryHeaderString) { - let millisToSleep = Math.round(parseFloat(retryHeaderString) * 1000) - if (isNaN(millisToSleep)) { - millisToSleep = Math.max(0, new Date(retryHeaderString) - new Date()) - } - return millisToSleep -} - -async function fetchAndRetryIfNecessary(callAPIFn) { - const response = await callAPIFn() - if (response.status === 429) { - const retryAfter = response.headers.get('retry-after') - const millisToSleep = getMillisToSleep(retryAfter) - await sleep(millisToSleep) - return fetchAndRetryIfNecessary(callAPIFn) - } - return response -} - -// Look up a Webex person by email via /v1/people. Used by the admin UI so -// authorizing someone on a bot only requires their email — the id, display -// name, and avatar come from Webex. Resolves to null when the email doesn't -// match any account or when the Webex request fails; the caller distinguishes -// with its own 404/502 as appropriate. -function findPersonByEmail(email) { - return new Promise(function (resolve, reject) { - var url = 'https://webexapis.com/v1/people?email=' + encodeURIComponent(email); - var requestOptions = { - method: 'GET', - headers: { - 'Authorization': 'Bearer ' + getServiceAccountAccessToken(), - 'Content-Type': 'application/json' - } - }; - fetchWithRateLimit(url, requestOptions) - .then(function (response) { - if (!response.ok) { - return reject(new Error('HTTP ' + response.status + ' from Webex /people')); - } - return response.json(); - }) - .then(function (body) { - var items = body && body.items; - if (!items || !items.length) return resolve(null); - var person = items[0]; - resolve({ - id: person.id, - displayName: person.displayName, - email: (person.emails && person.emails[0]) || email, - avatar: person.avatar || null - }); - }) - .catch(reject); - }); -} - -// Best-effort DELETE on a Webex message we previously sent (right now used -// only to retract the preview DM when the user cancels a draft). Never -// throws: the caller doesn't want a Webex hiccup here to bounce the whole -// cancel flow, so we swallow + log and let the caller move on. -function deleteWebexMessage(messageId, appName) { - return new Promise(function (resolve) { - var url = "https://webexapis.com/v1/messages/" + encodeURIComponent(messageId); - var requestOptions = { - method: 'DELETE', - headers: { "Authorization": "Bearer " + getBotToken(appName) } - }; - fetchWithRateLimit(url, requestOptions) - .then(function (response) { - if (!response.ok) { - logger('deleteWebexMessage', - 'HTTP ' + response.status + ' for ' + messageId); - } - resolve(response.ok); - }) - .catch(function (err) { - logger('deleteWebexMessage', - 'error for ' + messageId + ': ' + (err && err.message || err)); - resolve(false); - }); - }); -} - -function sendDirectMessage(toPersonId, message, imageName, appName) { - return new Promise(async function (resolve, reject) { - - var startTime = new Date(); - var myHeaders = new Headers(); - myHeaders.append("Authorization", "Bearer " + getBotToken(appName)); - var formdata = new FormData(); - formdata.append("toPersonId", toPersonId); - formdata.append("markdown", message); - - if (imageName) { - //var file = await fs.readFileSync("./uploads/" + imageName); - var stream = await fs.createReadStream("./uploads/" + imageName); - //console.log(file); - formdata.append("files", stream, imageName); - //formdata.set("files", new BlobFromStream(stream, content.length), imageName) - - }; - - //var myBody = JSON.stringify(myBodyJSON); - var requestOptions = { - method: 'POST', - headers: myHeaders, - body: formdata, - redirect: 'follow' - }; - - fetchWithRateLimit("https://webexapis.com/v1/messages", requestOptions) - .then(response => response.text()) - .then(result => { - var msgResult = JSON.parse(result); - resolve(msgResult) - }) - .catch(error => reject(error)); - }); -} - -function sendDirectCard(toPersonId, cardInfo, appName) { - return new Promise(async function (resolve, reject) { - - var myHeaders = { - "Authorization": "Bearer " + getBotToken(appName), - "Content-Type": "application/json" - } - var body = { - toPersonId: toPersonId, - markdown: cardInfo.message, - attachments: [{ contentType: "application/vnd.microsoft.card.adaptive", content: cardInfo.card }] - } - - var requestOptions = { - method: 'POST', - headers: myHeaders, - body: JSON.stringify(body), - redirect: 'follow' - }; - //console.log("RequestOptions:" + JSON.stringify(requestOptions)); - fetchWithRateLimit("https://webexapis.com/v1/messages", requestOptions) - .then(async function (response) { - var subData = await response.text(); - //console.log(subData); - resolve(response); - }) - .catch(error => reject(error)); - - }); -} -function buildTranslations(message) { - return new Promise(async function (resolve, reject) { - logger("buildTranslations", "Started"); - var translationPromises = []; - var translations = { - "english": { - "text": message.text, - "markdown": message.markdown, - "html": message.html - } - }; - - for (var language of config.languages) { - translationPromises.push(translateMessage(message.text, language, "text", "text")) - translationPromises.push(translateMessage(message.markdown, language, "text", "markdown")) - translationPromises.push(translateMessage(message.html, language, "html", "html")) - } - - Promise.allSettled(translationPromises) - .then(results => { - //console.log(JSON.stringify(results)) - for (var translation of results) { - if (translation.status == "fulfilled") { - if (translations[translation.value.name]) { - translations[translation.value.name][translation.value.messageType] = translation.value.translatedText; - } else { - translations[translation.value.name] = {}; - translations[translation.value.name][translation.value.messageType] = translation.value.translatedText; - } - - } - } - //console.log(JSON.stringify(translations)); - logger("buildTranslations", "Completed") - resolve(translations); - }) - }) -} - -function translateMessage(message, language, format, messageType) { - return new Promise(async function (resolve, reject) { - var urlSearchParams = { - q: message, - target: language.key, - format: format, - source: "en", - model: "nmt", - key: process.env.GOOGLE_TRANSLATE_API_KEY || '' - } - var translateUrl = "https://translation.googleapis.com/language/translate/v2?" + new URLSearchParams(urlSearchParams); - //console.log(translateUrl); - fetchWithRateLimit(translateUrl, { "method": "POST" }) - .then(response => response.json()) - .then(result => { - var languageResult = { - "name": language.name, - "key": language.key, - "format": format, - "messageType": messageType, - "translatedText": result.data.translations[0].translatedText - } - resolve(languageResult); - }) - .catch(error => reject(error)) - }) -} -function refreshToken() { - return new Promise(async function (resolve, reject) { - var url = "https://webexapis.com/v1/access_token"; - - var params = new URLSearchParams(); - params.append('grant_type', 'refresh_token'); - params.append('client_id', process.env.WEBEX_SERVICE_ACCOUNT_CLIENT_ID || ''); - params.append('client_secret', process.env.WEBEX_SERVICE_ACCOUNT_CLIENT_SECRET || ''); - params.append('refresh_token', serviceAccountToken.refresh_token); - - fetchWithRateLimit(url, { method: 'post', body: params }) - .then(res => { - console.log(res.status); - return res.json(); - }) - .then(function (json) { - var updated = Object.assign({}, json, { - created: new Date().toLocaleString("en-US", { timeZoneName: "short", timeZone: "America/New_York" }), - renewBy: new Date(Date.now() + (json.expires_in * 1000)).toLocaleString("en-US", { timeZoneName: "short", timeZone: "America/New_York" }), - refreshBy: new Date(Date.now() + (json.refresh_token_expires_in * 1000)).toLocaleString("en-US", { timeZoneName: "short", timeZone: "America/New_York" }) - }); - saveServiceAccountToken(updated); - resolve("Refresh Success!\nExpires: " + updated.renewBy + "\nRefreshBy: " + updated.refreshBy); - }) - .catch(err => reject(err)); - }) -} - function saveConfig(jsonObject, configFile) { const jsonData = JSON.stringify(jsonObject, null, 4); fs.writeFileSync(configFile, jsonData, (err) => { @@ -1665,10 +955,6 @@ function isAdmin(personId) { return helpers.isAdmin(authorized, personId); } -function msToTime(duration) { - return helpers.msToTime(duration); -} - function logger(activeFunction, logLine) { var d = new Date(); console.log(d.toLocaleString() + " " + activeFunction + ": " + logLine); diff --git a/lib/jobs.js b/lib/jobs.js new file mode 100644 index 0000000..a6e73f0 --- /dev/null +++ b/lib/jobs.js @@ -0,0 +1,324 @@ +// 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, + }; +} diff --git a/lib/translation.js b/lib/translation.js new file mode 100644 index 0000000..276d405 --- /dev/null +++ b/lib/translation.js @@ -0,0 +1,84 @@ +// Google Translate wrappers used when a bot is configured to fan a message +// out into multiple languages. Splitting out into its own module keeps the +// google-translate API key (env-only) and the per-language plural-request +// dance in one place — index.js doesn't need to know about either. +// +// createTranslator({ languages, apiKey, fetchWithRateLimit, logger }) +// returns { buildTranslations(message) }. Callers hand in the same +// fetchWithRateLimit that lib/webex.js uses so the whole app has a single +// 429 retry policy. + +export function createTranslator(deps) { + var languages = Array.isArray(deps.languages) ? deps.languages : []; + var apiKey = deps.apiKey || ''; + var fetchWithRateLimit = deps.fetchWithRateLimit; + var logger = deps.logger || function () {}; + + // POST to Google's translate v2 REST endpoint for one (message, target + // language, format) triple. Resolves to a normalized record so the + // fan-out caller can index by language name. + function translateMessage(message, language, format, messageType) { + return new Promise(function (resolve, reject) { + var params = { + q: message, + target: language.key, + format: format, + source: 'en', + model: 'nmt', + key: apiKey + }; + var translateUrl = 'https://translation.googleapis.com/language/translate/v2?' + new URLSearchParams(params); + fetchWithRateLimit(translateUrl, { method: 'POST' }) + .then(response => response.json()) + .then(result => { + resolve({ + name: language.name, + key: language.key, + format: format, + messageType: messageType, + translatedText: result.data.translations[0].translatedText + }); + }) + .catch(reject); + }); + } + + // Fans `message` out across every configured language in text + + // markdown + html variants. Returns an object shaped like: + // { english: {text, markdown, html}, "": {...}, ... } + // Failed translations are silently dropped (partial success is better + // than the whole send blowing up when one language server is down). + function buildTranslations(message) { + return new Promise(function (resolve) { + logger('buildTranslations', 'Started'); + var translationPromises = []; + var translations = { + english: { + text: message.text, + markdown: message.markdown, + html: message.html + } + }; + + for (var language of languages) { + translationPromises.push(translateMessage(message.text, language, 'text', 'text')); + translationPromises.push(translateMessage(message.markdown, language, 'text', 'markdown')); + translationPromises.push(translateMessage(message.html, language, 'html', 'html')); + } + + Promise.allSettled(translationPromises).then(function (results) { + for (var translation of results) { + if (translation.status === 'fulfilled') { + var name = translation.value.name; + if (!translations[name]) translations[name] = {}; + translations[name][translation.value.messageType] = translation.value.translatedText; + } + } + logger('buildTranslations', 'Completed'); + resolve(translations); + }); + }); + } + + return { buildTranslations, translateMessage }; +} diff --git a/lib/webex.js b/lib/webex.js new file mode 100644 index 0000000..dd20a48 --- /dev/null +++ b/lib/webex.js @@ -0,0 +1,409 @@ +// Webex REST client wrappers. Every function that hits webexapis.com lives +// here so index.js doesn't have to juggle URLs, auth headers, or the 429 +// retry loop. State (token getters, logger) is injected via +// createWebexClient() so the module itself stays free of module-level +// mutable state — that keeps it importable from anywhere without +// initialization-order landmines. +// +// Usage: +// import { createWebexClient } from './lib/webex.js'; +// const webex = createWebexClient({ +// getBotToken: (app) => botTokens[app] && botTokens[app].token, +// getServiceAccountAccessToken: () => serviceAccountToken.access_token, +// getServiceAccountRefreshToken: () => serviceAccountToken.refresh_token, +// saveServiceAccountToken: (t) => { /* persist to disk */ }, +// env: process.env, +// logger, +// }); + +import fs from 'fs'; + +export function createWebexClient(deps) { + var getBotToken = deps.getBotToken; + var getServiceAccountAccessToken = deps.getServiceAccountAccessToken; + var getServiceAccountRefreshToken = deps.getServiceAccountRefreshToken; + var saveServiceAccountToken = deps.saveServiceAccountToken; + var env = deps.env || {}; + var logger = deps.logger || function () {}; + + // ---- low-level fetch helpers ----------------------------------------- + + async function fetchWithRateLimit(url, requestOptions) { + const response = await fetch(url, requestOptions); + if (response.status === 429) { + const secondsToWait = Number(response.headers.get('retry-after')); + logger('fetchWithRateLimit', + 'Waiting ' + secondsToWait + 's on 429 for ' + url); + await new Promise(resolve => setTimeout(resolve, secondsToWait * 1000)); + return await fetchWithRateLimit(url, requestOptions); + } + return response; + } + + function getMillisToSleep(retryHeaderString) { + let millisToSleep = Math.round(parseFloat(retryHeaderString) * 1000); + if (isNaN(millisToSleep)) { + millisToSleep = Math.max(0, new Date(retryHeaderString) - new Date()); + } + return millisToSleep; + } + + function sleep(milliseconds) { + return new Promise((resolve) => setTimeout(resolve, milliseconds)); + } + + // Variant of fetchWithRateLimit that takes a callable so the caller can + // rebuild a fresh request per retry (needed for POSTs whose body is a + // stream — the body is consumed on the first attempt). + async function fetchAndRetryIfNecessary(callAPIFn) { + const response = await callAPIFn(); + if (response.status === 429) { + const retryAfter = response.headers.get('retry-after'); + const millisToSleep = getMillisToSleep(retryAfter); + await sleep(millisToSleep); + return fetchAndRetryIfNecessary(callAPIFn); + } + return response; + } + + // ---- people ---------------------------------------------------------- + + // GET /people/me — resolves the caller identified by `bearerToken`. + // Used by loadBotProfile (with the bot token) and the OAuth callback + // (with the freshly issued user access token). + function whoAmI(bearerToken) { + return new Promise(function (resolve, reject) { + logger('whoAmI', 'Fetching /people/me.'); + var requestOptions = { + method: 'GET', + headers: { + 'Authorization': 'Bearer ' + bearerToken, + 'Content-Type': 'application/json' + }, + redirect: 'follow' + }; + fetchWithRateLimit('https://webexapis.com/v1/people/me', requestOptions) + .then(async function (response) { + resolve(await response.json()); + }) + .catch(reject); + }); + } + + // GET /people?email= via the service account. Resolves to the raw + // Webex person record (id, displayName, emails, avatar, ...) or + // rejects on !ok. Used by the CSV recipient parser. + function getPersonInfo(email) { + return new Promise(function (resolve, reject) { + var url = 'https://webexapis.com/v1/people?email=' + encodeURIComponent(email) + '&max=100'; + var requestOptions = { + method: 'GET', + headers: { + 'Authorization': 'Bearer ' + getServiceAccountAccessToken(), + 'Content-Type': 'application/json' + }, + redirect: 'follow' + }; + fetchWithRateLimit(url, requestOptions) + .then(async function (response) { + logger('getPersonInfo', email + ': ' + response.status + ' ' + response.statusText); + if (!response.ok) return reject(new Error('HTTP ' + response.status)); + var userInfo = await response.json(); + resolve(userInfo.items && userInfo.items[0]); + }) + .catch(reject); + }); + } + + // GET /people?email= via the service account, normalized for the admin + // UI (id/displayName/email/avatar or null when no match). Wraps + // getPersonInfo behavior in the shape the /admin/users route wants. + function findPersonByEmail(email) { + return new Promise(function (resolve, reject) { + var url = 'https://webexapis.com/v1/people?email=' + encodeURIComponent(email); + var requestOptions = { + method: 'GET', + headers: { + 'Authorization': 'Bearer ' + getServiceAccountAccessToken(), + 'Content-Type': 'application/json' + } + }; + fetchWithRateLimit(url, requestOptions) + .then(function (response) { + if (!response.ok) { + return reject(new Error('HTTP ' + response.status + ' from Webex /people')); + } + return response.json(); + }) + .then(function (body) { + var items = body && body.items; + if (!items || !items.length) return resolve(null); + var person = items[0]; + resolve({ + id: person.id, + displayName: person.displayName, + email: (person.emails && person.emails[0]) || email, + avatar: person.avatar || null + }); + }) + .catch(reject); + }); + } + + // ---- groups ---------------------------------------------------------- + + // Paginated GET /groups. Returns the flat list of every group in the + // org (typically ~10k). Uses the service account token because bot + // tokens can't read the SCIM group directory. + function findWebexGroup() { + return new Promise(async function (resolve, reject) { + var groups = []; + var startIndex = 1; + var count = 500; + var groupSize = 10000; + + var url = new URL('https://webexapis.com/v1/groups'); + url.searchParams.append('attributes', 'displayName'); + url.searchParams.append('sortBy', 'displayName'); + url.searchParams.append('sortOrder', 'ascending'); + url.searchParams.append('includeMembers', 'false'); + url.searchParams.append('startIndex', startIndex); + url.searchParams.append('count', count); + + var requestOptions = { + method: 'GET', + headers: { + 'Authorization': 'Bearer ' + getServiceAccountAccessToken(), + 'Content-Type': 'application/json' + } + }; + while (groupSize > groups.length) { + try { + var fetchedGroups = await fetchWithRateLimit(url, requestOptions); + if (fetchedGroups.ok) { + var groupData = await fetchedGroups.json(); + groupSize = groupData.totalResults; + for (var group of groupData.groups) groups.push(group); + startIndex = startIndex + count; + url.searchParams.delete('startIndex'); + url.searchParams.append('startIndex', startIndex); + } else { + return reject(fetchedGroups.status + ': ' + fetchedGroups.statusText); + } + } catch (error) { + logger('findWebexGroup', 'page fetch failed: ' + (error && error.message || error)); + return reject(error); + } + } + resolve(groups); + }); + } + + // Paginated GET /groups/:id/members. Returns { members: [...], groups: + // [...] } separated by whether the child is a person or a nested + // group. Callers must recurse into nested groups themselves. + function getGroupMembers(groupId) { + return new Promise(async function (resolve, reject) { + var groupMembers = { members: [], groups: [] }; + var startIndex = 1; + var count = 500; + var memberSize = 10000; + + var requestOptions = { + method: 'GET', + headers: { + 'Authorization': 'Bearer ' + getServiceAccountAccessToken(), + 'Content-Type': 'application/json' + } + }; + + while (memberSize > (groupMembers.members.length + groupMembers.groups.length)) { + var pageUrl = 'https://webexapis.com/v1/groups/' + groupId + '/members?startIndex=' + startIndex + '&count=' + count; + try { + var groupResponse = await fetchWithRateLimit(pageUrl, requestOptions); + if (groupResponse.ok) { + var groupData = await groupResponse.json(); + groupMembers.displayName = groupData.displayName; + memberSize = groupData.memberSize; + for (var member of groupData.members) { + if (member.type == 'group') { + groupMembers.groups.push({ id: member.id, name: member.displayName }); + } else { + groupMembers.members.push(member); + } + } + startIndex = startIndex + count; + } else { + return reject(groupResponse.status + ': ' + groupResponse.statusText); + } + } catch (error) { + logger('getGroupMembers', 'page fetch failed: ' + (error && error.message || error)); + return reject(error); + } + } + resolve(groupMembers); + }); + } + + // ---- messages -------------------------------------------------------- + + // POST /messages as a bot. Attaches an image from ./uploads/ + // when supplied. Resolves with the parsed Webex message record + // (including id) so callers can retract later via deleteMessage. + function sendDirectMessage(toPersonId, message, imageName, appName) { + return new Promise(function (resolve, reject) { + var myHeaders = new Headers(); + myHeaders.append('Authorization', 'Bearer ' + getBotToken(appName)); + var formdata = new FormData(); + formdata.append('toPersonId', toPersonId); + formdata.append('markdown', message); + + if (imageName) { + var stream = fs.createReadStream('./uploads/' + imageName); + formdata.append('files', stream, imageName); + } + + var requestOptions = { + method: 'POST', + headers: myHeaders, + body: formdata, + redirect: 'follow' + }; + fetchWithRateLimit('https://webexapis.com/v1/messages', requestOptions) + .then(response => response.text()) + .then(result => { + try { resolve(JSON.parse(result)); } + catch (err) { reject(err); } + }) + .catch(reject); + }); + } + + // POST /messages carrying an Adaptive Card attachment. Fire-and-forget + // for job-completion cards. + function sendDirectCard(toPersonId, cardInfo, appName) { + return new Promise(function (resolve, reject) { + var body = { + toPersonId: toPersonId, + markdown: cardInfo.message, + attachments: [ + { + contentType: 'application/vnd.microsoft.card.adaptive', + content: cardInfo.card + } + ] + }; + var requestOptions = { + method: 'POST', + headers: { + 'Authorization': 'Bearer ' + getBotToken(appName), + 'Content-Type': 'application/json' + }, + body: JSON.stringify(body), + redirect: 'follow' + }; + fetchWithRateLimit('https://webexapis.com/v1/messages', requestOptions) + .then(async function (response) { + await response.text(); + resolve(response); + }) + .catch(reject); + }); + } + + // POST /messages via the retry-friendly path (streamed body). Used by + // the send queue where transient 429s are worth retrying. + function sendMessageWithRetry(toPersonId, message, imageName, appName) { + var url = 'https://webexapis.com/v1/messages'; + return fetchAndRetryIfNecessary(async () => { + var myHeaders = new Headers(); + myHeaders.append('Authorization', 'Bearer ' + getBotToken(appName)); + var formdata = new FormData(); + formdata.append('toPersonId', toPersonId); + formdata.append('markdown', message); + if (imageName) { + var stream = fs.createReadStream('./uploads/' + imageName); + formdata.append('files', stream, imageName); + } + return fetch(url, { + method: 'POST', + headers: myHeaders, + body: formdata, + redirect: 'follow' + }); + }); + } + + // Best-effort DELETE /messages/:id. Never rejects — the callers (right + // now: preview-DM retraction) always want the flow to continue whether + // this succeeds or not. Resolves to true on 2xx, false otherwise. + function deleteMessage(messageId, appName) { + return new Promise(function (resolve) { + var url = 'https://webexapis.com/v1/messages/' + encodeURIComponent(messageId); + var requestOptions = { + method: 'DELETE', + headers: { 'Authorization': 'Bearer ' + getBotToken(appName) } + }; + fetchWithRateLimit(url, requestOptions) + .then(function (response) { + if (!response.ok) { + logger('deleteMessage', 'HTTP ' + response.status + ' for ' + messageId); + } + resolve(response.ok); + }) + .catch(function (err) { + logger('deleteMessage', 'error for ' + messageId + ': ' + (err && err.message || err)); + resolve(false); + }); + }); + } + + // ---- OAuth refresh --------------------------------------------------- + + // Renews the service-account access token using its refresh token. + // Also computes / persists renewBy + refreshBy timestamps in + // America/New_York for the cron that reads them. + function refreshToken() { + return new Promise(function (resolve, reject) { + var params = new URLSearchParams(); + params.append('grant_type', 'refresh_token'); + params.append('client_id', env.WEBEX_SERVICE_ACCOUNT_CLIENT_ID || ''); + params.append('client_secret', env.WEBEX_SERVICE_ACCOUNT_CLIENT_SECRET || ''); + params.append('refresh_token', getServiceAccountRefreshToken()); + + fetchWithRateLimit('https://webexapis.com/v1/access_token', { method: 'post', body: params }) + .then(function (res) { + logger('refreshToken', 'HTTP ' + res.status); + return res.json(); + }) + .then(function (json) { + var opts = { timeZoneName: 'short', timeZone: 'America/New_York' }; + var updated = Object.assign({}, json, { + created: new Date().toLocaleString('en-US', opts), + renewBy: new Date(Date.now() + (json.expires_in * 1000)).toLocaleString('en-US', opts), + refreshBy: new Date(Date.now() + (json.refresh_token_expires_in * 1000)).toLocaleString('en-US', opts) + }); + saveServiceAccountToken(updated); + resolve('Refresh Success!\nExpires: ' + updated.renewBy + '\nRefreshBy: ' + updated.refreshBy); + }) + .catch(reject); + }); + } + + return { + fetchWithRateLimit, + fetchAndRetryIfNecessary, + getMillisToSleep, + sleep, + whoAmI, + getPersonInfo, + findPersonByEmail, + findWebexGroup, + getGroupMembers, + sendDirectMessage, + sendDirectCard, + sendMessageWithRetry, + deleteMessage, + refreshToken, + }; +}