242 lines
7.7 KiB
JavaScript
242 lines
7.7 KiB
JavaScript
function buildLiveRevisionQuery(revisions = {}) {
|
|
const params = new URLSearchParams();
|
|
const tokenPattern = /^[0-9a-f]{16}\.[0-9]{1,20}$/;
|
|
['context', 'events', 'notifications'].forEach((section) => {
|
|
const revision = revisions[section];
|
|
if (typeof revision === 'string' && tokenPattern.test(revision)) {
|
|
params.set(section + '_revision', revision);
|
|
}
|
|
});
|
|
return params.toString();
|
|
}
|
|
|
|
function retryAfterMs(value) {
|
|
if (typeof value !== 'string' || !/^\d+$/.test(value)) return null;
|
|
const seconds = Number(value);
|
|
if (!Number.isFinite(seconds) || seconds <= 0) return null;
|
|
return Math.min(seconds * 1000, 300000);
|
|
}
|
|
|
|
function createContextPoller({
|
|
fetchContext,
|
|
onSnapshot,
|
|
onError,
|
|
isHidden = () => false,
|
|
setTimer = setTimeout,
|
|
clearTimer = clearTimeout,
|
|
setDeadlineTimer = setTimeout,
|
|
clearDeadlineTimer = clearTimeout,
|
|
intervalMs = 8000,
|
|
timeoutMs = 12000,
|
|
maxBackoffMs = 60000,
|
|
}) {
|
|
let activeRequest = null;
|
|
let timer = null;
|
|
let stopped = false;
|
|
let revisions = {};
|
|
let retainedSnapshot = null;
|
|
let nextDelayMs = intervalMs;
|
|
let failureStreak = 0;
|
|
let nextRetryAt = null;
|
|
let lastSuccessAt = null;
|
|
|
|
function snapshotDelay(snapshot) {
|
|
const freshness = snapshot && snapshot.freshness;
|
|
const sections = freshness && freshness.sections;
|
|
const sectionValues = sections && Object.values(sections);
|
|
const degradedRetrySeconds = (sectionValues || [])
|
|
.filter(section => section.degraded)
|
|
.map(section => Number(section.retry_in_seconds ?? freshness.retry_in_seconds))
|
|
.filter(seconds => Number.isFinite(seconds) && seconds > 0);
|
|
const failedSectionDelay = degradedRetrySeconds.length ?
|
|
Math.min(...degradedRetrySeconds) * 1000 : null;
|
|
if (sectionValues && sectionValues.length && sectionValues.every(section => section.degraded)) {
|
|
return failedSectionDelay || intervalMs;
|
|
}
|
|
const freshForSeconds = Number(freshness && freshness.fresh_for_seconds);
|
|
const healthyDeadlines = (sectionValues || [])
|
|
.filter(section => !section.degraded)
|
|
.map(section => freshForSeconds - Number(section.age_seconds))
|
|
.filter(seconds => Number.isFinite(seconds));
|
|
if (Number.isFinite(freshForSeconds) && freshForSeconds > 0 && healthyDeadlines.length) {
|
|
const earliestDeadline = Math.min(...healthyDeadlines);
|
|
const healthyDelay = earliestDeadline <= 0 &&
|
|
sectionValues.some(section => !section.degraded && section.revalidating) ?
|
|
intervalMs : Math.max(0, earliestDeadline) * 1000;
|
|
return failedSectionDelay === null ? healthyDelay : Math.min(healthyDelay, failedSectionDelay);
|
|
}
|
|
return intervalMs;
|
|
}
|
|
|
|
function cancelTimer() {
|
|
if (timer !== null) clearTimer(timer);
|
|
timer = null;
|
|
}
|
|
|
|
function schedule(delayMs = intervalMs) {
|
|
cancelTimer();
|
|
if (stopped || isHidden()) return;
|
|
nextRetryAt = Date.now() + delayMs;
|
|
timer = setTimer(() => {
|
|
timer = null;
|
|
nextRetryAt = null;
|
|
refresh();
|
|
}, delayMs);
|
|
}
|
|
|
|
function abortError() {
|
|
const error = new Error('Live update superseded.');
|
|
error.name = 'AbortError';
|
|
return error;
|
|
}
|
|
|
|
function supersede(request) {
|
|
if (!request || request.settled) return;
|
|
request.superseded = true;
|
|
request.controller.abort();
|
|
request.rejectDeadline(abortError());
|
|
}
|
|
|
|
function refresh(options = {}) {
|
|
if (stopped || isHidden()) return Promise.resolve(null);
|
|
if (activeRequest && !options.force) return activeRequest.promise;
|
|
if (activeRequest) supersede(activeRequest);
|
|
|
|
const controller = new AbortController();
|
|
let rejectDeadline;
|
|
const deadlinePromise = new Promise((resolve, reject) => {
|
|
rejectDeadline = reject;
|
|
});
|
|
const requestState = {
|
|
controller,
|
|
deadline: null,
|
|
promise: null,
|
|
rejectDeadline,
|
|
settled: false,
|
|
superseded: false,
|
|
};
|
|
activeRequest = requestState;
|
|
requestState.deadline = setDeadlineTimer(() => {
|
|
if (requestState.settled || requestState.superseded) return;
|
|
controller.abort();
|
|
const error = new Error(`Live update timed out after ${timeoutMs}ms.`);
|
|
error.name = 'TimeoutError';
|
|
rejectDeadline(error);
|
|
}, timeoutMs);
|
|
|
|
let request;
|
|
try {
|
|
request = fetchContext(options.full ? {} : { ...revisions }, { signal: controller.signal });
|
|
} catch (error) {
|
|
request = Promise.reject(error);
|
|
}
|
|
requestState.promise = Promise.race([Promise.resolve(request), deadlinePromise])
|
|
.then((snapshot) => {
|
|
if (activeRequest !== requestState) return retainedSnapshot;
|
|
failureStreak = 0;
|
|
lastSuccessAt = Date.now();
|
|
nextDelayMs = snapshotDelay(snapshot);
|
|
const changedSections = ['context', 'events', 'notifications'].filter(
|
|
(section) => Object.prototype.hasOwnProperty.call(snapshot, section)
|
|
);
|
|
retainedSnapshot = retainedSnapshot ? { ...retainedSnapshot, ...snapshot } : { ...snapshot };
|
|
revisions = { ...revisions, ...(snapshot.revisions || {}) };
|
|
onSnapshot(retainedSnapshot, changedSections);
|
|
return retainedSnapshot;
|
|
})
|
|
.catch((error) => {
|
|
if (activeRequest === requestState && !stopped) {
|
|
const retryAfterMs = Number(error && error.retryAfterMs);
|
|
failureStreak += 1;
|
|
nextDelayMs = Number.isFinite(retryAfterMs) && retryAfterMs > 0
|
|
? retryAfterMs
|
|
: Math.min(intervalMs * (2 ** (failureStreak - 1)), maxBackoffMs);
|
|
onError(error);
|
|
}
|
|
return null;
|
|
})
|
|
.finally(() => {
|
|
requestState.settled = true;
|
|
clearDeadlineTimer(requestState.deadline);
|
|
if (activeRequest !== requestState) return;
|
|
activeRequest = null;
|
|
schedule(nextDelayMs);
|
|
});
|
|
return requestState.promise;
|
|
}
|
|
|
|
function setVisible(visible) {
|
|
cancelTimer();
|
|
if (!visible) return Promise.resolve(null);
|
|
return refresh({ force: true });
|
|
}
|
|
|
|
function adopt(snapshot) {
|
|
if (
|
|
stopped || activeRequest || retainedSnapshot || !snapshot ||
|
|
typeof snapshot !== 'object' ||
|
|
!Object.prototype.hasOwnProperty.call(snapshot, 'context')
|
|
) return false;
|
|
retainedSnapshot = { ...snapshot };
|
|
revisions = { ...(snapshot.revisions || {}) };
|
|
failureStreak = 0;
|
|
lastSuccessAt = Date.now();
|
|
nextDelayMs = snapshotDelay(snapshot);
|
|
const changedSections = ['context', 'events', 'notifications'].filter(
|
|
section => Object.prototype.hasOwnProperty.call(snapshot, section)
|
|
);
|
|
onSnapshot(retainedSnapshot, changedSections);
|
|
schedule(nextDelayMs);
|
|
return true;
|
|
}
|
|
|
|
async function adoptPending(snapshotPromise) {
|
|
if (!snapshotPromise || typeof snapshotPromise.then !== 'function') return false;
|
|
let deadline = null;
|
|
const timeout = new Promise(resolve => {
|
|
deadline = setDeadlineTimer(() => resolve(null), timeoutMs);
|
|
});
|
|
try {
|
|
const snapshot = await Promise.race([
|
|
Promise.resolve(snapshotPromise).catch(() => null),
|
|
timeout,
|
|
]);
|
|
return adopt(snapshot);
|
|
} finally {
|
|
clearDeadlineTimer(deadline);
|
|
}
|
|
}
|
|
|
|
return {
|
|
start: refresh,
|
|
refresh,
|
|
adopt,
|
|
adoptPending,
|
|
setVisible,
|
|
getState() {
|
|
return {
|
|
refreshing: Boolean(activeRequest),
|
|
nextRetryAt,
|
|
lastSuccessAt,
|
|
failureStreak,
|
|
};
|
|
},
|
|
stop() {
|
|
stopped = true;
|
|
cancelTimer();
|
|
if (activeRequest) {
|
|
const request = activeRequest;
|
|
activeRequest = null;
|
|
supersede(request);
|
|
}
|
|
},
|
|
};
|
|
}
|
|
|
|
createContextPoller.buildRevisionQuery = buildLiveRevisionQuery;
|
|
createContextPoller.retryAfterMs = retryAfterMs;
|
|
|
|
if (typeof module !== 'undefined' && module.exports) {
|
|
module.exports = createContextPoller;
|
|
}
|