#!/usr/bin/env node import fs from 'fs'; import fetch from 'node-fetch'; import { Headers } from 'node-fetch'; import express from 'express'; import bodyParser from 'body-parser'; import multer from 'multer'; import sharp from 'sharp'; import { v4 as uuid } from 'uuid'; import cookieParser from 'cookie-parser'; import FormData from 'form-data'; import cron from "node-cron"; import PQueue from 'p-queue'; //import pLimit from "p-limit"; //const limit = pLimit(10); const queue = new PQueue({ concurrency: 10 }); //Load the config files from storage. var config = JSON.parse(fs.readFileSync('./config/config.json')); var jobs = JSON.parse(fs.readFileSync('./config/jobs.json')); var userPrefs = JSON.parse(fs.readFileSync('./config/userPrefs.json')); // Defensive init so code below can safely `.push`, `.filter`, and index into // these regardless of what the on-disk jobs.json happens to contain. jobs.building = jobs.building || {}; jobs.running = Array.isArray(jobs.running) ? jobs.running : []; jobs.scheduled = Array.isArray(jobs.scheduled) ? jobs.scheduled : []; jobs.completed = Array.isArray(jobs.completed) ? jobs.completed : []; if (typeof jobs.lastJobNumber !== 'number') jobs.lastJobNumber = 0; // Per-bot tokens live in their own gitignored file so adding a new bot is one // JSON entry + an env-free restart. Shape: // { "": { "token": "...", "enabled": true } } var botTokens = loadBotTokens(); function loadBotTokens() { try { return JSON.parse(fs.readFileSync('./config/botTokens.json')); } catch (err) { console.log('loadBotTokens failed: ' + err.message); return {}; } } function getBotToken(appName) { var entry = botTokens[appName]; if (!entry) return null; if (typeof entry === 'string') return entry; // tolerate flat shape if (entry.enabled === false) return null; return entry.token || null; } function isBotEnabled(appName) { return getBotToken(appName) !== null; } // Returns the bot's config block if the bot is both defined in config.json AND // has an enabled token entry. Returns null otherwise. Use this everywhere // instead of reaching into config.webex.bot[...] directly so unknown or // disabled bots can't crash the request. function getBotConfig(appName) { if (!appName) return null; var cfg = config.webex && config.webex.bot && config.webex.bot[appName]; if (!cfg) return null; if (!isBotEnabled(appName)) return null; return cfg; } // Cached bot profiles populated from Webex /people/me at startup so the // front-end and the completion card can use the bot's actual displayName and // avatar without hardcoding either in config.json. // botProfiles[appName] = { displayName, avatar } var botProfiles = {}; async function loadBotProfile(appName) { var token = getBotToken(appName); if (!token) return; try { var response = await fetchWithRateLimit('https://webexapis.com/v1/people/me', { method: 'GET', headers: { 'Authorization': 'Bearer ' + token } }); if (!response.ok) { logger('loadBotProfile', appName + ' /people/me failed: ' + response.status + ' ' + response.statusText); return; } var data = await response.json(); botProfiles[appName] = { displayName: data.displayName, avatar: data.avatar, personId: data.id, emails: data.emails }; logger('loadBotProfile', appName + ' = ' + data.displayName); } catch (err) { logger('loadBotProfile', appName + ' error: ' + (err && err.message || err)); } } function loadBotProfiles() { var apps = Object.keys((config.webex && config.webex.bot) || {}); return Promise.allSettled(apps.filter(isBotEnabled).map(loadBotProfile)); } function getBotProfile(appName) { return botProfiles[appName] || null; } // Express middleware: 404s any /CollabCentral/:app/* request whose `:app` is // not a known + enabled bot. Attaches the bot config to req.botConfig for // downstream handlers. function requireBot(req, res, next) { var botCfg = getBotConfig(req.params.app); if (!botCfg) { logger("requireBot", "Unknown or disabled bot: " + req.params.app); return res.status(404).send("Unknown bot: " + req.params.app); } req.botConfig = botCfg; next(); } // Building (draft) jobs are keyed by cookieId AND appName so a user who is // authorized for more than one bot can have one independent draft per bot. function buildingKey(req) { return (req.cookies && req.cookies.id) + ':' + req.params.app; } // Filters one of the global job arrays (running/scheduled/completed) down to // the jobs that belong to the requested bot. function jobsForApp(arr, appName) { return (arr || []).filter(function (j) { return j && j.appName === appName; }); } // Service-account OAuth tokens are rewritten in place by refreshToken(), so // they live in their own file rather than mixed into config.json. var serviceAccountToken = loadServiceAccountToken(); function loadServiceAccountToken() { try { return JSON.parse(fs.readFileSync('./config/token.json')); } catch (err) { console.log('loadServiceAccountToken failed: ' + err.message); return {}; } } function saveServiceAccountToken(tokenObj) { // Keep the in-memory copy in sync even if the disk write fails so the // refreshed token is still usable for the rest of the process lifetime. serviceAccountToken = tokenObj; try { fs.writeFileSync('./config/token.json', JSON.stringify(tokenObj, null, 4)); } catch (err) { logger("saveServiceAccountToken", "Failed to persist token.json: " + err.message); } } function getServiceAccountAccessToken() { return serviceAccountToken && serviceAccountToken.access_token; } // Builds the Webex OAuth authorize URL for a given bot. The redirect_uri must // match one of the URIs registered with the Webex integration in the Webex // developer portal (one URI per bot). Returns null if required env vars are // missing so callers can return an actionable error instead of a malformed URL. function buildAuthUrl(appName) { var clientId = process.env.WEBEX_INTEGRATION_CLIENT_ID; var template = process.env.OAUTH_CALLBACK_URL_TEMPLATE; if (!clientId || !template) { logger('buildAuthUrl', 'Missing WEBEX_INTEGRATION_CLIENT_ID or OAUTH_CALLBACK_URL_TEMPLATE env var.'); return null; } var params = new URLSearchParams(); params.append('client_id', clientId); params.append('response_type', 'code'); params.append('redirect_uri', template.replace(':app', appName)); params.append('scope', 'spark:kms spark:people_read'); params.append('state', ''); return 'https://webexapis.com/v1/authorize?' + params.toString(); } function getOAuthRedirectUri(appName) { var template = process.env.OAUTH_CALLBACK_URL_TEMPLATE || ''; return template.replace(':app', appName); } //Load Express Server var app = express(); app.use(bodyParser.json()); app.use(cookieParser()); var serverPort = process.env.SERVER_PORT || config.server.port; var server = app.listen(serverPort, function () { logger("startup", config.server.name + " running on port " + serverPort + "."); loadBotProfiles() .then(() => logger("startup", "Loaded " + Object.keys(botProfiles).length + " bot profile(s).")) .catch(err => logger("startup", "loadBotProfiles error: " + (err && err.message || err))); }); server.setTimeout(3000000); const upload = multer({ limits: { fileSize: 4000000 } }).fields( [ { name: 'uploadImage' }, { name: 'uploadCSV' } ] ); // Anchor both crons to a specific timezone so the schedule is deterministic // regardless of the host/container clock. Override with CRON_TIMEZONE if the // deployment moves to a different region. var cronTimezone = process.env.CRON_TIMEZONE || 'America/New_York'; // Daily at 01:10 local time: prune completed jobs older than the retention // window and rewrite jobs.json. cron.schedule('0 10 1 * * *', async function () { var jobFile = JSON.parse(fs.readFileSync('./config/jobs.json')); var cleanJobs = await cleanCompletedJobs(jobFile); saveConfig(cleanJobs, './config/jobs.json'); jobs = JSON.parse(fs.readFileSync('./config/jobs.json')); }, { timezone: cronTimezone }) // Every minute: refresh the service-account token if it's inside the renewal // window, and dispatch any scheduled jobs whose time has arrived. cron.schedule('0 * * * * *', () => { var renewBy = serviceAccountToken && serviceAccountToken.renewBy; if (renewBy && new Date(new Date(renewBy) - 7200000) < new Date()) { logger("refreshToken", "Renew by: " + new Date(renewBy).toLocaleString()) logger("refreshToken", "Renew after: " + new Date(new Date(renewBy) - 7200000).toLocaleString()); logger("refreshToken", "Refreshing token.") refreshToken() .then(response => { logger("tokenRefresh", response); }) .catch(error => logger("refreshToken", error)); } checkScheduledJobs() .catch(error => logger("checkScheduledJobs", error)); }, { timezone: cronTimezone }) //Routes to be used // Validate the :app segment for every /CollabCentral/:app/* request. This // runs before the static mount and all the per-route handlers below, so any // unknown or disabled bot gets a clean 404 instead of crashing on undefined. app.use('/CollabCentral/:app', requireBot) app.use('/CollabCentral/:app', express.static('html')) app.get('/status', function (req, res) { res.status(200).send({ "status": "I'm alive!'" }); }); app.get('/CollabCentral/:app/authUrl', function (req, res) { var url = buildAuthUrl(req.params.app); if (!url) { return res.status(500).send('OAuth is not configured. Check WEBEX_INTEGRATION_CLIENT_ID and OAUTH_CALLBACK_URL_TEMPLATE.'); } res.send(url); }); // Returns the bot's display metadata + whether the current cookie is authorized. // The front-end uses this to drive page labels, icons, and the // "you don't have access" view, so adding a new bot doesn't require any code // changes in the HTML/JS. app.get('/CollabCentral/:app/info', function (req, res) { var appName = req.params.app; var botCfg = req.botConfig; // populated by requireBot var profile = getBotProfile(appName) || {}; var iconBase = botCfg.iconBase || appName; var personId = req.cookies && req.cookies.id; res.status(200).send({ appName: appName, label: botCfg.label || appName, displayName: profile.displayName || botCfg.label || appName, iconUrl: '/CollabCentral/' + appName + '/' + iconBase + '.png', faviconUrl: '/CollabCentral/' + appName + '/' + iconBase + '.ico', avatarUrl: profile.avatar || null, authorized: isAuthorized(appName, personId) }); }); app.post('/CollabCentral/:app/jobs/:action', (req, res) => { if (isAuthorized(req.params.app, req.cookies.id)) { req.setTimeout(3000000); if (req.params.action == "edit") { logger("apiEndpoint(" + req.params.app + ")", "POST /jobs/edit"); buildJob(req, res) .then(job => { 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 } await buildTranslations(jobs.building[buildingKey(req)].message.english) .then(result => { jobs.building[buildingKey(req)].message = result; }) .catch(error => console.log("Error translating: " + error)); res.status(200).send(job) saveConfig(jobs, "./config/jobs.json") }) .catch(error => console.log("Error sending message: " + error)); }) .catch(error => console.log("Error buildJob: " + error)); } else 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)]; processRunningQueue() .catch(error => console.log("sendRunningMessages error: " + error)) res.status(204).redirect("/CollabCentral/" + req.params.app + "/monitorJobs.html"); } else 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') res.status(204).redirect("/CollabCentral/" + req.params.app + "/monitorJobs.html"); } else { res.status(404) } } else { res.status(401).send("You are not authorized.") } }) app.get('/CollabCentral/:app/jobs/list/:scope', (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: jobsForApp(jobs.running, appName), scheduled: jobsForApp(jobs.scheduled, appName), completed: 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 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 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(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', (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); }); app.get('/CollabCentral/:app/user/:scope/:action', (req, res) => { console.log("Got /user/:scope/:action request."); //console.log(req.cookies); if (isAuthorized(req.params.app, req.cookies.id)) { //console.log("Person is authorized for the request.") if (req.params.scope == "groups") { if (req.params.action == "list") { logger("apiEndpoint(" + req.params.app + ")", "GET /user/groups/list"); res.status(200).send(config.webex.bot[req.params.app].authorized[req.cookies.id].groups) } else if (req.params.action == "find") { logger("apiEndpoint(" + req.params.app + ")", "GET /user/groups/find"); findWebexGroup() .then(response => { //console.log(response) saveConfig(response, "./testData.json") res.status(200).send(response); }) } } else { res.status(404) } } else { res.status(401) } }) // Standard options for the session cookies set after a successful OAuth round-trip. const SESSION_COOKIE_OPTIONS = { httpOnly: true, secure: true, sameSite: 'strict', maxAge: 24 * 60 * 60 * 1000 }; app.get(`/CollabCentral/:app/oauth`, async function (req, res) { var appName = req.params.app; var authCode = req.query.code; if (!process.env.WEBEX_INTEGRATION_CLIENT_ID || !process.env.WEBEX_INTEGRATION_CLIENT_SECRET || !process.env.OAUTH_CALLBACK_URL_TEMPLATE) { logger("oauth", "Missing Webex integration env vars; cannot complete OAuth."); return res.status(500).send("OAuth is not configured on the server."); } if (!authCode) { return res.status(400).send("Missing OAuth `code` query parameter."); } var url = new URL("https://webexapis.com/v1/access_token"); var params = new URLSearchParams(); params.append('grant_type', 'authorization_code'); params.append('client_id', process.env.WEBEX_INTEGRATION_CLIENT_ID); params.append('client_secret', process.env.WEBEX_INTEGRATION_CLIENT_SECRET); params.append('code', authCode); params.append('redirect_uri', getOAuthRedirectUri(appName)); var response; try { response = await 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."); } if (!response.ok) { var errText = await response.text().catch(() => ''); logger("oauth", "Token exchange returned " + response.status + " " + response.statusText + " " + errText); return res.status(401).send(response.statusText || "OAuth token exchange failed."); } var jsonData = await response.json(); var whoami; try { whoami = await 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."); } logger("oauth", whoami.displayName + " successfully authed for " + appName + "."); res .cookie('displayName', whoami.displayName, SESSION_COOKIE_OPTIONS) .cookie('id', whoami.id, SESSION_COOKIE_OPTIONS) .cookie('avatar', whoami.avatar, SESSION_COOKIE_OPTIONS) .cookie('email', whoami.userName, SESSION_COOKIE_OPTIONS) .cookie('orgId', whoami.orgId, SESSION_COOKIE_OPTIONS) .cookie('access_token', jsonData.access_token, SESSION_COOKIE_OPTIONS) .cookie('refresh_token', jsonData.refresh_token, SESSION_COOKIE_OPTIONS) .redirect(301, '/CollabCentral/' + appName + '/sendMessage.html'); }); function buildJob(req, res) { return new Promise(async function (resolve, reject) { res.setTimeout(3000000) var memberList = []; var memberListPromises = []; upload(req, res, async function (err) { //console.log(req.files); //console.log(req.body); 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": [] } // check for error thrown by multer- file size etc if (err || req.files === undefined) { //no file } else { if (req.files["uploadImage"]) { for (var file of req.files["uploadImage"]) { let fileName = uuid() + ".jpeg" var image = await sharp(file.buffer).jpeg({ quality: 40, }).toFile('./uploads/' + fileName).catch(err => { console.log('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(buildPeopleList(fileName)); } } } if (req.body.groups) { logger("buildJob", "Locating groups: " + req.body.groups); memberListPromises.push(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(peoplePromises => { console.log(JSON.stringify(peoplePromises)) for (var peopleList of peoplePromises) { if (peopleList.status == "fulfilled") { for (var person of peopleList.value) { memberList.push(person); } } } //Make user members unique var uniqueMembers = [...new Map(memberList.map((m) => [m.id, m])).values()]; jobs.building[buildingKey(req)].memberList = uniqueMembers; var numSending = uniqueMembers.length.toString(); saveConfig(jobs, "./config/jobs.json") resolve(jobs.building[buildingKey(req)]); }) }) }) } 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` } }; console.log(url) while (groupSize > groups.length) { console.log("groupSize: " + groupSize) console.log("groups.length: " + groups.length) try { console.log(url); var fetchedGroups = await fetchWithRateLimit(url, requestOptions) if (fetchedGroups.ok) { var groupData = await fetchedGroups.json(); console.log("Total Results: " + groupData.totalResults) 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 { memberSize = 0; reject(fetchedGroups.status + ": " + fetchedGroups.statusText); } } catch (error) { console.log("Error findGroupById: " + error); } } console.log("Groups found: " + groups.length); resolve(groups); }) } function whoAmI(bearerToken) { return new Promise(async function (resolve, reject) { logger("whoAmI", "Person is having an existential crisis of identify.") 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 } 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) => { if (err) { throw err; } else { logger("saveConfig", "Wrote " + configFile); } }) } // Completed jobs are kept for review for COMPLETED_RETENTION_DAYS. Anything // older is dropped from jobs.completed by the daily cron. Bumping this number // just means we hold more history (and a bigger jobs.json) on disk. const COMPLETED_RETENTION_DAYS = 30; async function cleanCompletedJobs(jobs) { var cutoffDate = new Date(Date.now() - (COMPLETED_RETENTION_DAYS * 24 * 60 * 60 * 1000)); if (!Array.isArray(jobs.completed)) { console.log('No completed array found – nothing to do.'); return jobs; } const originalCount = jobs.completed.length; jobs.completed = jobs.completed.filter(job => { // Prefer endTime, fall back to startTime if missing const jobDateStr = job.endTime || job.startTime || job.created; if (!jobDateStr) return true; // safety: keep if no date const jobDate = new Date(jobDateStr); return jobDate >= cutoffDate; }); const removed = originalCount - jobs.completed.length; console.log(`Removed ${removed} job(s) older than ${COMPLETED_RETENTION_DAYS} days (cutoff ${cutoffDate.toISOString().split('T')[0]}).`); console.log(`Remaining completed jobs: ${jobs.completed.length}`); return jobs; } function isAuthorized(appName, personId) { logger("isAuthorized", appName + " " + personId) var botCfg = getBotConfig(appName); if (!botCfg) return false; if (!personId) return false; return !!(botCfg.authorized && botCfg.authorized[personId]); } function msToTime(duration) { var milliseconds = parseInt((duration % 1000) / 100) , seconds = parseInt((duration / 1000) % 60) , minutes = parseInt((duration / (1000 * 60)) % 60) , hours = parseInt((duration / (1000 * 60 * 60)) % 24); //hours = (hours < 10) ? "0" + hours : hours; //minutes = (minutes < 10) ? "0" + minutes : minutes; //seconds = (seconds < 10) ? "0" + seconds : seconds; var resultTime = ""; if (hours > 0) { resultTime = hours + "h " + minutes + "m " + seconds + "." + milliseconds + "s"; } else if (minutes > 0) { resultTime = minutes + "m " + seconds + "." + milliseconds + "s"; } else { resultTime = seconds + "." + milliseconds + "s"; } return resultTime; } function logger(activeFunction, logLine) { var d = new Date(); console.log(d.toLocaleString() + " " + activeFunction + ": " + logLine); } // gracefully shutdown (ctrl-c) process.on('SIGINT', function () { server.close(() => { logger("shutdown", config.server.name + ' stopped!') process.exit(); }); })