Merge retry hardening with DSP cache opt-in

This commit is contained in:
Jonathan Sykes
2026-10-10 23:00:32 +08:00
10 changed files with 282 additions and 26 deletions

View File

@@ -60,12 +60,14 @@ async function call(zeroName, tauriName, payload = {}) {
// Generic fetch wrapper — returns parsed JSON or throws with a human message.
async function webFetch(path, opts = {}) {
if ((opts.method || 'GET').toUpperCase() === 'GET' && window.AsyncGuard?.getJson)
return AsyncGuard.getJson(path, opts, {
// Some upstream failures are represented by HTTP 200 + {ok:false}; do
// not expose those transient responses as a manual Retry state either.
retryResult: value => value?.ok === false && /\b(?:408|425|429|5\d\d)\b|temporar|timed?\s*out|network|connection|upstream|unavailable|available right now/i.test(String(value.error || value.message || '')),
});
const res = await fetch(path, opts);
if (!res.ok) {
let msg = `HTTP ${res.status}`;
try { const j = await res.json(); msg = j.error || msg; } catch { /* non-JSON */ }
throw new Error(msg);
}
if (!res.ok) throw new Error(`HTTP ${res.status}`);
return res.json();
}
@@ -573,7 +575,7 @@ const API = {
? webFetch(`/api/search?q=${encodeURIComponent(query)}${refresh ? '&refresh=1' : ''}`)
: call('yt.search', 'yt_search', { query }),
related: (meta, { refresh = false } = {}) => WEB && /^[\w-]{11}$/.test(meta.id)
? fetch(`/api/related?${new URLSearchParams({ videoId: meta.id, title: meta.title || '', channel: meta.channel || meta.artist || '', ...(refresh ? { refresh: '1' } : {}) })}`).then(response => response.json())
? webFetch(`/api/related?${new URLSearchParams({ videoId: meta.id, title: meta.title || '', channel: meta.channel || meta.artist || '', ...(refresh ? { refresh: '1' } : {}) })}`)
: API.search(meta.title || meta.channel, { refresh }),
// Videos the server already knows — answers fast while the real search loads.
searchLocal: (query) => WEB
@@ -8666,9 +8668,7 @@ const RUNNING_BUILD = (() => {
let _knownBuildTag = null;
async function fetchServerBuild() {
const res = await fetch('/api/version', { cache: 'no-store' });
if (!res.ok) return null;
const data = await res.json();
const data = await webFetch('/api/version', { cache: 'no-store' });
return data.buildTag || null;
}

View File

@@ -162,7 +162,7 @@
const completionKey = manifest => '/__ytp_completion/' + manifest.buildTag;
async function complete(manifest, { cache, fetchFn, cycles = 4, now = Date.now, optIns,
sleep = ms => new Promise(resolve => setTimeout(resolve, ms)), baseDelay = 1000,
maxDelay = 30000, notify = () => {}, force = false } = {}) {
maxDelay = 30000, notify = () => {}, force = false, concurrency = 1 } = {}) {
const key = completionKey(manifest), saved = await cache.match(key);
let job = saved ? await saved.json() : { failures: 0, nextRetryAt: 0 };
const selectedOptIns = optIns || await readOptIns(cache);
@@ -177,7 +177,7 @@
for (let cycle = 0; cycle < cycles; cycle++) {
await publish(true);
let error;
try { await syncAssets(manifest, { cache, fetchFn, all: true, optIns: selectedOptIns }); }
try { await syncAssets(manifest, { cache, fetchFn, all: true, optIns: selectedOptIns, concurrency }); }
catch (failure) { error = failure; }
snapshot = await completeness(manifest, cache, { optIns: selectedOptIns });
if (snapshot.offlineReady) { job.failures = 0; job.nextRetryAt = 0; job.error = null; await publish(false); return job; }

View File

@@ -96,6 +96,13 @@ test('overlapping layout and completion jobs share a six-download ceiling',async
assert.ok(peak<=6);assert.ok([...calls.values()].every(count=>count===1));
});
test('offline completion keeps its own network concurrency to one',async()=>{
const c=cache(),m=manifest();let active=0,peak=0;
for(let i=0;i<12;i++){const p='/complete'+i+'.js';m.files[p]={h:String(i)};}
const fetchFn=async u=>{active++;peak=Math.max(peak,active);await new Promise(r=>setTimeout(r,2));active--;return response(u.split('=')[1]);};
await core.complete(m,{cache:c,fetchFn});assert.equal(peak,1);
});
test('blocking commit never waits on an optional file readiness probe',{timeout:1000},async()=>{
const c=cache(),m=manifest();m.files['/extra.js']={h:'extra'};m.groups.extra={files:['/extra.js'],background:true,contract:1};
await c.put('/a.js?v=a',response('a'));await c.put('/index.html?v=b',response('b'));

View File

@@ -35,7 +35,68 @@
}
}
const AsyncGuard = { runExclusive };
const wait = ms => new Promise(resolve => setTimeout(resolve, ms));
async function retry(task, { attempts = 3, delay = attempt => attempt === 1 ? 350 : 1000,
shouldRetry = () => true, sleep = wait } = {}) {
let last;
for (let attempt = 1; attempt <= Math.max(1, attempts); attempt++) {
try { return await task(attempt); }
catch (error) {
last = error;
if (attempt >= attempts || !shouldRetry(error, attempt)) throw error;
await sleep(delay(attempt));
}
}
throw last;
}
// Retry only idempotent GETs. Keep the timeout active through JSON parsing so
// a stalled response body is treated like a stalled connection.
async function getJson(url, options = {}, config = {}) {
const fetchFn = config.fetch || root.fetch?.bind(root);
if (!fetchFn) throw new Error('fetch is unavailable');
if ((options.method || 'GET').toUpperCase() !== 'GET') throw new Error('Retry helper only accepts GET');
const attempts = config.attempts ?? 3, timeoutMs = config.timeoutMs ?? 20000;
return retry(async () => {
const controller = typeof AbortController === 'function' ? new AbortController() : null;
let timedOut = false, timer, externalAbort;
const external = options.signal;
if (controller && external) {
if (external.aborted) throw external.reason || new DOMException('Aborted', 'AbortError');
externalAbort = () => controller.abort(external.reason);
external.addEventListener('abort', externalAbort, { once: true });
}
if (controller && timeoutMs > 0) timer = setTimeout(() => { timedOut = true; controller.abort(); }, timeoutMs);
try {
const response = await fetchFn(url, { ...options, ...(controller ? { signal: controller.signal } : {}) });
if ([408, 425, 429].includes(response.status) || response.status >= 500) {
const error = new Error(`HTTP ${response.status}`); error.retryable = true; throw error;
}
if (!response.ok) {
let message = `HTTP ${response.status}`;
try { message = (await response.json()).error || message; } catch {}
const error = new Error(message); error.retryable = false; throw error;
}
const value = await response.json();
if (config.retryResult?.(value)) {
const error = new Error(value.error || value.message || 'Temporary API failure');
error.retryable = true;
throw error;
}
return value;
} catch (error) {
if (external?.aborted) throw error;
if (timedOut) { const timeout = new Error('Request timed out'); timeout.retryable = true; throw timeout; }
if (error.retryable === undefined) error.retryable = true;
throw error;
} finally {
if (timer) clearTimeout(timer);
if (external && externalAbort) external.removeEventListener('abort', externalAbort);
}
}, { attempts, shouldRetry: error => error.retryable !== false, sleep: config.sleep || wait });
}
const AsyncGuard = { runExclusive, retry, getJson };
if (typeof module !== 'undefined' && module.exports) {
module.exports = AsyncGuard;

View File

@@ -60,6 +60,55 @@ test('frees the key after failure so a retry is possible', async () => {
assert.strictEqual(out, 'ok');
});
test('bounded GET retries recover transient network and server failures', async () => {
const { getJson } = require('./async-guard');
let calls = 0;
const result = await getJson('/api/search', {}, { sleep: async () => {}, fetch: async () => {
calls++;
if (calls === 1) throw new TypeError('connection reset');
if (calls === 2) return new Response('{}', { status: 503 });
return Response.json({ ok: true, results: [] });
} });
assert.equal(result.ok, true);
assert.equal(calls, 3);
});
test('GET retry stops at the bound and does not retry permanent HTTP errors', async () => {
const { getJson } = require('./async-guard');
let calls = 0;
await assert.rejects(getJson('/api/search', {}, { sleep: async () => {}, fetch: async () => {
calls++; return new Response('{}', { status: 404 });
} }), /HTTP 404/);
assert.equal(calls, 1);
await assert.rejects(getJson('/api/search', {}, { attempts: 2, sleep: async () => {}, fetch: async () => {
calls++; throw new TypeError('offline');
} }), /offline/);
assert.equal(calls, 3);
});
test('retryable API error payloads recover but permanent application errors do not retry', async () => {
const { getJson } = require('./async-guard'); let calls = 0;
const retryResult = value => value?.ok === false && /temporarily unavailable/i.test(value.error || '');
const result = await getJson('/api/streams', {}, { retryResult, sleep: async () => {}, fetch: async () => {
calls++;
return Response.json(calls === 1 ? { ok: false, error: 'temporarily unavailable' } : { ok: true, data: [] });
} });
assert.equal(result.ok, true); assert.equal(calls, 2);
calls = 0;
const absent = await getJson('/api/streams', {}, { retryResult, sleep: async () => {}, fetch: async () => { calls++; return Response.json({ ok: false, error: 'video not found' }); } });
assert.equal(absent.ok, false); assert.equal(calls, 1);
});
test('timed out GET bodies retry within the configured bound', async () => {
const { getJson } = require('./async-guard'); let calls = 0;
const result = await getJson('/api/related', {}, { attempts: 2, timeoutMs: 5, sleep: async () => {}, fetch: async (_url, { signal }) => {
calls++;
if (calls === 1) return new Promise((resolve, reject) => signal.addEventListener('abort', () => reject(new DOMException('aborted', 'AbortError')), { once: true }));
return Response.json({ ok: true });
} });
assert.equal(result.ok, true); assert.equal(calls, 2);
});
test('supports a synchronous fn', async () => {
const set = new Set();
const out = await runExclusive(set, 's', () => 7);

View File

@@ -6,6 +6,7 @@
try { manifest = JSON.parse(doc.getElementById('ytp-assets').textContent); }
catch { return; } // The index must supply this map; never fetch one during boot.
const assets = new Map(), groups = new Map(), hooks = new Map();
let warmScheduled = false;
const css = new Set([...doc.querySelectorAll('link[data-lazy-css]')].map(link => link.dataset.lazyCss));
const selected = new Map();
const url = path => selected.get(path) || (manifest.files[path]?.h ? path + '?v=' + manifest.files[path].h : './' + path.replace(/^\//,''));
@@ -50,14 +51,18 @@
const style = /\.css$/.test(path);
if (style && css.has(path)) return Promise.resolve();
if (!style && !/\.js$/.test(path)) return Promise.resolve(); // workers/imports/fonts are not page scripts
const task = new Promise((resolve, reject) => {
const loadOnce = () => new Promise((resolve, reject) => {
const node = doc.createElement(style ? 'link' : 'script');
const timer = setTimeout(() => { node.remove?.(); reject(new Error('Timed out loading ' + path)); }, 20000);
if (style) { node.rel = 'stylesheet'; node.href = url(path); }
else { node.async = false; node.src = url(path); }
node.onload = resolve;
node.onerror = () => reject(new Error('Unable to load ' + path));
node.onload = () => { clearTimeout(timer); resolve(); };
node.onerror = () => { clearTimeout(timer); node.remove?.(); reject(new Error('Unable to load ' + path)); };
doc.head.append(node);
});
const task = root.AsyncGuard?.retry
? root.AsyncGuard.retry(loadOnce, { attempts: 3 })
: loadOnce();
assets.set(path, task);
task.catch(() => assets.delete(path));
return task;
@@ -100,10 +105,17 @@
else root.navigator.serviceWorker?.getRegistration().then(send).catch(()=>{});
},
warm() {
if (!root.navigator.serviceWorker || root.navigator.onLine === false) return;
root.navigator.serviceWorker.ready.then(reg => {
for (const worker of new Set([reg.active, reg.waiting])) worker?.postMessage({type:'COMPLETE_ASSETS'});
}).catch(()=>{});
if (warmScheduled || !root.navigator.serviceWorker || root.navigator.onLine === false) return;
warmScheduled = true;
const start = () => {
warmScheduled = false;
if(root.navigator.onLine === false)return;
root.navigator.serviceWorker.ready.then(reg => {
for (const worker of new Set([reg.active, reg.waiting])) worker?.postMessage({type:'COMPLETE_ASSETS'});
}).catch(()=>{});
};
if (root.requestIdleCallback) root.requestIdleCallback(start, { timeout: 8000 });
else setTimeout(start, 5000);
},
onLoad(name, fn) {
if(api.loaded(name)) { Promise.resolve().then(fn).catch(error=>{if(root.CustomEvent)root.dispatchEvent?.(new root.CustomEvent('ytp-lazy-error',{detail:{group:name,error}}));}); return; }

View File

@@ -7,7 +7,7 @@ function fixture(current, held = [], appCache = false, userAgent = 'AppleWebKit/
const manifest = { buildTag: 'own', appCache, groups: { core: { files: [] }, 'layout:classic': { files: ['/classic.css'] }, 'layout:glass-stage': { files: ['/glass.css', '/one.js', '/two.js'] }, 'feature:test': { contract:1, files: ['/one.js', '/two.js'] } }, files: { '/app.js':{h:'a'}, '/classic.css': {h:'c'}, '/glass.css':{h:'g'}, '/one.js':{h:'1'}, '/two.js':{h:'2'} } };
const doc = { readyState:'loading', documentElement: { dataset:{} }, getElementById: () => ({ textContent:JSON.stringify(manifest) }), querySelectorAll: () => [], createElement: tag => ({ tagName:tag.toUpperCase() }), write: value => writes.push(value), addEventListener: (t,f) => listeners[t]=f };
doc.head = { append: node => { inserted.push(node); queueMicrotask(() => node.onload?.()); } };
const root = { document:doc, localStorage:{ getItem:() => '{"settings":{"layout":"glass-stage"}}' }, console, Promise, URL, crypto:{subtle:{}}, fetch:async key=>{fetched.push(key);return new Response('app');}, setTimeout, clearTimeout, navigator:{userAgent,serviceWorker:{ready:Promise.resolve({}),controller:controlled?{}:null,addEventListener(){}}}, addEventListener(){} };
const root = { document:doc, localStorage:{ getItem:() => '{"settings":{"layout":"glass-stage"}}' }, console, Promise, URL, crypto:{subtle:{}}, AsyncGuard:require('./async-guard'), requestIdleCallback:fn=>{root.__idle=fn;}, fetch:async key=>{fetched.push(key);return new Response('app');}, setTimeout, clearTimeout, navigator:{userAgent,serviceWorker:{ready:Promise.resolve({}),controller:controlled?{}:null,addEventListener(){}}}, addEventListener(){} };
if (current) root.caches = { open: async () => ({ match: async key => key === '/__ytp_asset_state' ? new Response(JSON.stringify({current})) : held.includes(key) ? new Response('', {headers:{'X-Asset-Hash':key.split('=')[1]}}) : undefined }) };
root.window=root; root.globalThis=root;
vm.runInNewContext(readFileSync(require.resolve('./lazy.js'),'utf8'), root);
@@ -27,11 +27,13 @@ test('lazy groups execute ordered classic scripts once, sharing concurrent calls
await root.Lazy.load('layout:glass-stage');
assert.equal(inserted.filter(n=>n.tagName==='SCRIPT').length,2);
});
test('failed injection can retry and proxy dispatches only after execution', async () => {
test('transient script injection failure retries before proxy dispatch', async () => {
const {root,inserted}=fixture();
const append=root.document.head.append;
root.document.head.append=n=>{inserted.push(n);queueMicrotask(()=>n.onerror?.());};
await assert.rejects(root.Lazy.load('feature:test'),/one.js/);
let first=true;
root.document.head.append=n=>{inserted.push(n);queueMicrotask(()=>{if(first){first=false;n.onerror?.();}else n.onload?.();});};
await root.Lazy.load('feature:test');
assert.equal(inserted.length,3,'first file is retried, then the second file loads');
root.document.head.append=append;
root.TestApi={open:value=>value+1};
const proxy=root.Lazy.proxy('feature:test',['open'],'TestApi');
@@ -92,6 +94,8 @@ test('online launches request worker completion even with Save-Data and skip off
f.root.navigator.connection={saveData:true};f.root.navigator.onLine=true;
f.root.navigator.serviceWorker={ready:Promise.resolve({active,waiting})};
f.root.Lazy.warm();await Promise.resolve();
assert.equal(sent.length,0,'warming waits until the page is idle');
f.root.__idle();await Promise.resolve();
assert.deepEqual(sent.map(m=>m.type),['COMPLETE_ASSETS','COMPLETE_ASSETS']);
f.root.navigator.onLine=false;f.root.Lazy.warm();await Promise.resolve();assert.equal(sent.length,2);
});

View File

@@ -265,7 +265,8 @@ self.addEventListener('fetch', (e) => {
return;
}
if (ASSET_SYNC && request.mode === 'navigate') e.waitUntil(completeAssets());
// The page schedules completion after first paint/idle via Lazy.warm(). Do
// not start a six-file background burst from this foreground fetch event.
// Only intercept GET/HEAD — let POST (sync endpoint) go through unmodified
if (request.method !== 'GET' && request.method !== 'HEAD') return;
@@ -597,8 +598,8 @@ async function activateAssets() {
for (const req of await cache.keys()) if(new URL(req.url).pathname.startsWith('/__ytp_candidate/')) await cache.delete(req);
await self.clients.claim();
if(old && old.current.buildTag !== m.buildTag) for(const c of await self.clients.matchAll({type:'window',includeUncontrolled:true})) c.postMessage({type:'SW_UPDATE_AVAILABLE',version:VERSION});
// Completion owns a message-event lifetime, not the activation barrier.
// In particular, a legacy client's four-second reload must not wait for it.
// Resume completion after activation, but with one request at a time so
// foreground navigation retains most of the available bandwidth.
self.registration.active.postMessage({ type: 'COMPLETE_ASSETS' });
try { console.info('[asset-sync] storage', await self.navigator.storage.estimate()); } catch {}
}