collabcentral/index.js
Joseph B. McQueen e604c7e9c9 Initial commit: multi-bot CollabCentral
Extends the single-Novi codebase into a multi-bot mass-messenger where
each bot has its own token, avatar, label, and per-user authorization.

- Secrets moved out of config.json: per-bot tokens in gitignored
  config/botTokens.json (with enabled flag), service-account OAuth in
  gitignored config/token.json (rewritten by refresh cron), integration
  and Google keys in .env.
- Single Webex integration handles OAuth for all bots via a per-app
  redirect URI derived from OAUTH_CALLBACK_URL_TEMPLATE.
- New requireBot middleware and getBotConfig helper reject requests for
  unknown or disabled bots at the /CollabCentral/:app boundary.
- New /info endpoint plus dynamic frontend loading (sendMessage,
  monitorJobs) so pages self-describe per bot, including bot avatar
  fetched from Webex /people/me at startup.
- Job draft state keyed by cookieId + appName so each bot has its own
  building queue; job list/detail endpoints filter by appName so users
  only see jobs from bots they are authorized on.
- New jobDetail page for a readable per-job view; completed jobs are
  retained for 30 days by the cleanup cron.
- Completion adaptive cards use per-bot avatar and label.
- Miscellaneous fixes: off-by-two in the send loop, removed three dead
  send/process variants, added defensive init for jobs.* on load,
  dropped the deprecated crypto npm shim, and cleaned up stray logger
  labels and typos.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-01 17:53:07 -04:00

1346 lines
53 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/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:
// { "<appName>": { "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) {
fs.writeFileSync('./config/token.json', JSON.stringify(tokenObj, null, 4));
serviceAccountToken = tokenObj;
}
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' }
]
);
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'));
})
cron.schedule('0 * * * * *', () => {
var startRefreshTime = new Date();
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 => console.log("refreshToken: " + error));
}
checkScheduledJobs()
.then((result) => {
console.log(result)
})
.catch(error => logIt("checkScheduledJobs: " + error));
})
//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);
var count = 0;
queue.on('active', () => {
console.log(`Working on item #${++count}. Size: ` + queue.size + ` Pending: ` + queue.pending);
});
queue.on('completed', async function (result) {
console.log(JSON.stringify(result));
});
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 => {
console.log(JSON.stringify.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();
});
})