collabcentral/lib/webex.js
Joseph B. McQueen 8c3a32f47d Phase R2: extract lib/webex, lib/translation, lib/jobs from index.js
index.js had grown to 1,642 lines / ~50 functions. This peels the
three biggest self-contained concerns out into their own modules
and wires them back in through a dependency-bag factory so each
module stays free of module-level mutable state.

lib/webex.js — every webexapis.com round-trip (fetchWithRateLimit,
whoAmI, findWebexGroup, getGroupMembers, getPersonInfo,
findPersonByEmail, sendDirectMessage, sendMessageWithRetry,
sendDirectCard, deleteMessage, refreshToken). Factory closes over
token getters and the logger.

lib/translation.js — Google Translate fan-out (buildTranslations,
translateMessage). Reuses the webex 429 helper so the app has one
retry policy.

lib/jobs.js — cron-driven pipeline (checkScheduledJobs,
processRunningQueue, sendQueueJobMessages, buildJobCompletedCard)
plus the two recipient-resolution helpers (collectGroupMembers,
buildPeopleList). Factory takes jobs/queue/userPrefs/webex/etc so
mutable state stays owned by index.js.

index.js: instantiates webex/translator/jobsPipeline once, rewires
every call site to go through them, and drops ~715 lines of moved
code plus a dead msToTime wrapper. Down from 1,642 → 969 lines,
26 top-level functions instead of 50+.

Tests: 53/53 helpers still green. Smoke tested /info, /admin/users,
/user/groups/list, /user/groups/find, /jobs/list/all, and the
unauth 401 path — all match pre-R2 behavior.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-02 08:33:42 -04:00

409 lines
17 KiB
JavaScript

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