fix: recover live refresh from stalled requests (#303)
All checks were successful
CI / lint (pull_request) Successful in 31s
CI / build-frontend (pull_request) Successful in 5s

This commit is contained in:
timmy 2026-08-08 13:15:20 +00:00
parent 7356530318
commit 40c66df143
5 changed files with 193 additions and 14 deletions

View File

@ -5,9 +5,12 @@ function createContextPoller({
isHidden = () => false,
setTimer = setTimeout,
clearTimer = clearTimeout,
setDeadlineTimer = setTimeout,
clearDeadlineTimer = clearTimeout,
intervalMs = 8000,
timeoutMs = 12000,
}) {
let inFlight = null;
let activeRequest = null;
let timer = null;
let stopped = false;
let revisions = {};
@ -27,18 +30,55 @@ function createContextPoller({
}, intervalMs);
}
function refresh() {
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 (inFlight) return inFlight;
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({ ...revisions });
request = fetchContext({ ...revisions }, { signal: controller.signal });
} catch (error) {
request = Promise.reject(error);
}
inFlight = Promise.resolve(request)
requestState.promise = Promise.race([Promise.resolve(request), deadlinePromise])
.then((snapshot) => {
if (activeRequest !== requestState) return retainedSnapshot;
const changedSections = ['context', 'events', 'notifications'].filter(
(section) => Object.prototype.hasOwnProperty.call(snapshot, section)
);
@ -48,20 +88,23 @@ function createContextPoller({
return retainedSnapshot;
})
.catch((error) => {
onError(error);
if (activeRequest === requestState && !stopped) onError(error);
return null;
})
.finally(() => {
inFlight = null;
requestState.settled = true;
clearDeadlineTimer(requestState.deadline);
if (activeRequest !== requestState) return;
activeRequest = null;
schedule();
});
return inFlight;
return requestState.promise;
}
function setVisible(visible) {
cancelTimer();
if (!visible) return Promise.resolve(null);
return refresh();
return refresh({ force: true });
}
return {
@ -71,6 +114,11 @@ function createContextPoller({
stop() {
stopped = true;
cancelTimer();
if (activeRequest) {
const request = activeRequest;
activeRequest = null;
supersede(request);
}
},
};
}

View File

@ -229,7 +229,7 @@
function setClock() { qs('#clock').textContent = fmt(new Date()); }
setClock(); setInterval(setClock, 1000);
async function fetchLiveSnapshot(revisions = {}) {
async function fetchLiveSnapshot(revisions = {}, { signal } = {}) {
const params = new URLSearchParams();
Object.entries(revisions).forEach(([section, revision]) => {
if (Number.isInteger(revision) && revision >= 0) params.set(section + '_revision', revision);
@ -237,6 +237,7 @@
const query = params.toString();
const res = await fetch('api/v1/live' + (query ? '?' + query : ''), {
headers: { Accept: 'application/json' },
signal,
});
if (!res.ok) throw new Error('HTTP ' + res.status);
return res.json();
@ -608,7 +609,9 @@
liveMode = false;
activeFlushLogin = '';
if (!hasContextSnapshot && hydrateOfflineWork('outage')) return;
setStatus(hasContextSnapshot ? 'Update failed · showing last snapshot' : 'Unavailable');
const timeoutStatus = e.name === 'TimeoutError' ?
'Update delayed · showing last snapshot' : 'Update failed · showing last snapshot';
setStatus(hasContextSnapshot ? timeoutStatus : 'Unavailable');
if (!hasContextSnapshot) {
qs('#context').innerHTML = '<div class="muted">Context unavailable.</div>';
qs('#view-hint').textContent = 'Active view unavailable.';
@ -2811,7 +2814,7 @@
isHidden: () => document.hidden,
intervalMs: 8000,
});
function load() { return contextPoller.refresh(); }
function load() { return contextPoller.refresh({ force: true }); }
const offlineStatus = qs('#offline-status');
const keepWorkOffline = qs('#keep-work-offline');
@ -2884,7 +2887,7 @@
offlineStatus.hidden = true;
setOfflineWorkMode(false);
setStatus('Reconnecting…');
contextPoller.refresh();
contextPoller.refresh({ force: true });
}
keepWorkOffline.addEventListener('change', () => {
offlineWorkStore.setEnabled(keepWorkOffline.checked);

View File

@ -114,3 +114,113 @@ const poller = createContextPoller({{
{"context": "timmy", "event": 2, "age": 2, "changed": []},
],
}
def test_context_poller_aborts_a_stalled_request_and_recovers_on_schedule():
script = f"""
const createContextPoller = require({json.dumps(str(POLLER))});
const timers = [];
const deadlineTimers = [];
const errors = [];
const snapshots = [];
const signals = [];
let calls = 0;
const poller = createContextPoller({{
fetchContext: (revisions, options) => {{
calls += 1;
signals.push(options.signal);
if (calls === 1) return new Promise(() => {{}});
return Promise.resolve({{ context: {{ user: {{ login: 'timmy' }} }} }});
}},
onSnapshot: snapshot => snapshots.push(snapshot.context.user.login),
onError: error => errors.push({{ name: error.name, message: error.message }}),
setTimer: (callback, delay) => {{
const timer = {{ callback, delay, cancelled: false }};
timers.push(timer);
return timer;
}},
clearTimer: timer => {{ timer.cancelled = true; }},
setDeadlineTimer: (callback, delay) => {{
const timer = {{ callback, delay, cancelled: false }};
deadlineTimers.push(timer);
return timer;
}},
clearDeadlineTimer: timer => {{ timer.cancelled = true; }},
intervalMs: 8000,
timeoutMs: 25,
}});
(async () => {{
const first = poller.start();
const deadline = deadlineTimers.find(timer => timer.delay === 25 && !timer.cancelled);
deadline.callback();
await first;
const retry = timers.find(timer => timer.delay === 8000 && !timer.cancelled);
retry.callback();
await new Promise(resolve => setImmediate(resolve));
process.stdout.write(JSON.stringify({{
calls,
firstAborted: signals[0].aborted,
errors,
snapshots,
}}));
}})();
"""
assert run_node(script) == {
"calls": 2,
"firstAborted": True,
"errors": [
{"name": "TimeoutError", "message": "Live update timed out after 25ms."}
],
"snapshots": ["timmy"],
}
def test_forced_refresh_retires_stale_generation_without_late_overwrite():
script = f"""
const createContextPoller = require({json.dumps(str(POLLER))});
let resolveFirst;
let firstSignal;
const rendered = [];
const errors = [];
let calls = 0;
const poller = createContextPoller({{
fetchContext: (revisions, options) => {{
calls += 1;
if (calls === 1) {{
firstSignal = options.signal;
return new Promise(resolve => {{ resolveFirst = resolve; }});
}}
return Promise.resolve({{ context: {{ user: {{ login: 'fresh' }} }} }});
}},
onSnapshot: snapshot => rendered.push(snapshot.context.user.login),
onError: error => errors.push(error.name),
setTimer: () => 1,
clearTimer: () => {{}},
setDeadlineTimer: () => 1,
clearDeadlineTimer: () => {{}},
}});
(async () => {{
const stale = poller.start();
const fresh = poller.refresh({{ force: true }});
await fresh;
resolveFirst({{ context: {{ user: {{ login: 'stale' }} }} }});
await stale;
await new Promise(resolve => setImmediate(resolve));
process.stdout.write(JSON.stringify({{
calls,
firstAborted: firstSignal.aborted,
rendered,
errors,
}}));
}})();
"""
assert run_node(script) == {
"calls": 2,
"firstAborted": True,
"rendered": ["fresh"],
"errors": [],
}

View File

@ -57,3 +57,21 @@ async def test_dashboard_reports_each_live_section_from_its_own_freshness():
assert "const notificationsFresh = Array.isArray(snapshot.notifications)" not in html
assert "Unread updates unavailable · showing last known updates" in html
assert "Activity refresh failed · showing last activity" in html
@pytest.mark.anyio
async def test_operator_refresh_and_reconnect_supersede_a_stalled_live_request():
html = await dashboard()
assert "async function fetchLiveSnapshot(revisions = {}, { signal } = {})" in html
assert "headers: { Accept: 'application/json' },\n signal," in html
assert "function load() { return contextPoller.refresh({ force: true }); }" in html
assert "contextPoller.refresh({ force: true });" in html
@pytest.mark.anyio
async def test_stalled_refresh_keeps_the_last_snapshot_with_specific_guidance():
html = await dashboard()
assert "e.name === 'TimeoutError'" in html
assert "Update delayed · showing last snapshot" in html

View File

@ -103,7 +103,7 @@ async def test_dashboard_announces_offline_mode_and_refreshes_after_reconnect():
bootstrap = (Path(__file__).resolve().parents[1] / "frontend" / "dashboard.js").read_text()
assert "window.addEventListener('offline'" in bootstrap
assert "window.addEventListener('online'" in bootstrap
assert "contextPoller.refresh()" in bootstrap
assert "contextPoller.refresh({ force: true })" in bootstrap
@pytest.mark.anyio