From 0251720fa201a194755758b77dfe7e19f73c7014 Mon Sep 17 00:00:00 2001 From: Jonathan Sykes Date: Fri, 9 Oct 2026 09:32:08 +0800 Subject: [PATCH] Complete verified offline assets in a resumable worker job --- frontend/asset-sync-core.js | 84 +++++++++++++++++++++++++++++--- frontend/asset-sync-core.test.js | 36 ++++++++++++++ frontend/asset-worker.test.js | 9 ++-- frontend/sw.js | 43 ++++++++++++---- 4 files changed, 150 insertions(+), 22 deletions(-) diff --git a/frontend/asset-sync-core.js b/frontend/asset-sync-core.js index e8d951c..3efa777 100644 --- a/frontend/asset-sync-core.js +++ b/frontend/asset-sync-core.js @@ -6,7 +6,7 @@ const url = (path, file) => path + '?v=' + file.h; function paths(manifest) { - return [...new Set(Object.values(manifest.groups).flatMap(group => group.files))].sort(); + return Object.keys(manifest.files).sort(); } function blocking(manifest, previous, activeLayout) { @@ -55,13 +55,33 @@ return !!response && response.headers.get('X-Asset-Hash') === manifest.files[path].h; })); const missing = found.filter(value => !value).length; - return { ready: missing === 0, missing, version: manifest.buildTag }; + const offline = await completeness(manifest, cache); + return { ready: missing === 0, missing, version: manifest.buildTag, ...offline }; } - async function syncAssets(manifest, { cache, fetchFn, concurrency = 6, attempts = 3, activeLayout, groups }) { + // All concurrent messages/layout jobs in this worker share six download slots. + let downloads = 0; + const waiters = []; + async function downloadSlot(task) { + if (downloads >= 6) await new Promise(resolve => waiters.push(resolve)); + else downloads++; + try { return await task(); } + finally { const next = waiters.shift(); if (next) next(); else downloads--; } + } + + const pendingDownloads = new Map(); + async function sharedDownload(key, task) { + while (pendingDownloads.has(key)) await pendingDownloads.get(key).catch(() => {}); + const promise = downloadSlot(task); + pendingDownloads.set(key, promise); + try { return await promise; } + finally { pendingDownloads.delete(key); } + } + + async function syncAssets(manifest, { cache, fetchFn, concurrency = 6, attempts = 3, activeLayout, groups, all = false }) { const previous = (await state(cache))?.current; const missing = []; - const selected = groups ? [...new Set(groups.flatMap(name => manifest.groups[name]?.files || []))].sort() + const selected = all ? paths(manifest) : groups ? [...new Set(groups.flatMap(name => manifest.groups[name]?.files || []))].sort() : blocking(manifest, previous, activeLayout); for (const path of selected) { const key = url(path, manifest.files[path]); @@ -71,8 +91,11 @@ const count = missing.length; let failure; await Promise.all(Array.from({ length: Math.max(1, Math.min(6, concurrency, count)) }, async () => { - while (missing.length && !failure) { + while (missing.length) { const key = missing.shift(); + try { await sharedDownload(key, async () => { + const held = await cache.match(key); + if (held?.headers.get('X-Asset-Hash') === key.split('v=')[1]) return; let error; for (let attempt = 0; attempt < Math.min(3, attempts); attempt++) { const controller = new AbortController(); @@ -89,14 +112,59 @@ } catch (err) { error = err; } finally { clearTimeout(timer); } } - if (error) failure = error; + if (error) throw error; + }); } catch (error) { failure = error; } } })); if (failure) throw failure; - if (!groups && !(await status(manifest, cache, { previous, activeLayout })).ready) throw new Error('Incomplete blocking assets'); + if (!groups && !all && !(await status(manifest, cache, { previous, activeLayout })).ready) throw new Error('Incomplete blocking assets'); return { refreshed: count, caches: 1 }; } + + // Verified cache entries are the durable progress journal. Counters are rebuilt + // after worker termination or browser eviction; never trust a stored ready bit. + async function completeness(manifest, cache) { + const missingFiles = []; + for (const path of paths(manifest)) { + const response = await cache.match(url(path, manifest.files[path])); + if (!response || response.headers.get('X-Asset-Hash') !== manifest.files[path].h) missingFiles.push(path); + } + const total = paths(manifest).length; + return { offlineReady: missingFiles.length === 0, cached: total - missingFiles.length, total, missingFiles }; + } + + const completionKey = manifest => '/__ytp_completion/' + manifest.buildTag; + async function complete(manifest, { cache, fetchFn, cycles = 4, now = Date.now, + sleep = ms => new Promise(resolve => setTimeout(resolve, ms)), baseDelay = 1000, + maxDelay = 30000, notify = () => {}, force = false } = {}) { + const key = completionKey(manifest), saved = await cache.match(key); + let job = saved ? await saved.json() : { failures: 0, nextRetryAt: 0 }; + let snapshot = await completeness(manifest, cache); + const publish = async running => { + job = { ...job, ...snapshot, running, version: manifest.buildTag }; + await cache.put(key, Response.json(job)); + await notify(job); + }; + if (snapshot.offlineReady) { job.failures = 0; job.nextRetryAt = 0; await publish(false); return job; } + if (!force && job.nextRetryAt > now()) await sleep(Math.min(maxDelay, job.nextRetryAt - now())); + for (let cycle = 0; cycle < cycles; cycle++) { + await publish(true); + let error; + try { await syncAssets(manifest, { cache, fetchFn, all: true }); } + catch (failure) { error = failure; } + snapshot = await completeness(manifest, cache); + if (snapshot.offlineReady) { job.failures = 0; job.nextRetryAt = 0; job.error = null; await publish(false); return job; } + job.failures++; + const delay = Math.min(maxDelay, baseDelay * 2 ** Math.min(job.failures - 1, 10)); + job.nextRetryAt = now() + delay; + job.error = String(error?.message || 'Incomplete offline cache'); + await publish(cycle + 1 < cycles); + if (cycle + 1 < cycles) await sleep(delay); + } + return job; + } + async function commit(manifest, cache, { previousClients, activeLayout } = {}) { const old = await state(cache); if (!(await status(manifest, cache, { previous: old?.current, activeLayout })).ready) { @@ -119,7 +187,7 @@ return next; } - const api = { CACHE, STATE, url, paths, blocking, plan, retained, fallback, status, syncAssets, state, commit }; + const api = { CACHE, STATE, url, paths, blocking, plan, retained, fallback, status, completeness, completionKey, complete, syncAssets, state, commit }; if (typeof module !== 'undefined' && module.exports) module.exports = api; else root.AssetSyncCore = api; })(typeof globalThis !== 'undefined' ? globalThis : this); diff --git a/frontend/asset-sync-core.test.js b/frontend/asset-sync-core.test.js index 2c9a9d4..7eb64ec 100644 --- a/frontend/asset-sync-core.test.js +++ b/frontend/asset-sync-core.test.js @@ -45,3 +45,39 @@ test('removing group membership retains the old tab URL even if the physical fil const previous=manifest(),current=manifest();current.groups.core.files=['/index.html']; assert.ok(core.retained(current,previous).has('/a.js?v=a')); }); + +test('completion includes ungrouped runtime files and exposes verified missing paths', async () => { + const c=cache(), m=manifest();m.files['/worker.js']={h:'w'}; + await core.syncAssets(m,{cache:c,fetchFn:async u=>response(u.split('=')[1])}); + const status=await core.status(m,c); + assert.equal(status.ready,true);assert.equal(status.offlineReady,false); + assert.deepEqual(status.missingFiles,['/worker.js']);assert.equal(status.cached,2);assert.equal(status.total,3); + await core.complete(m,{cache:c,fetchFn:async u=>response(u.split('=')[1])}); + assert.equal((await core.status(m,c)).offlineReady,true); +}); + +test('completion persists backoff and resumes only missing files after a worker restart',async()=>{ + const c=cache(),m=manifest();let clock=1000;const calls=[]; + await core.complete(m,{cache:c,cycles:1,now:()=>clock,baseDelay:10,fetchFn:async u=>{if(u.includes('index'))throw Error('offline');return response('a');}}); + const checkpoint=await(await c.match(core.completionKey(m))).json(); + assert.equal(checkpoint.running,false);assert.equal(checkpoint.failures,1);assert.equal(checkpoint.nextRetryAt,1010); + await core.complete(m,{cache:c,now:()=>clock,sleep:async ms=>{assert.equal(ms,10);clock+=ms;},fetchFn:async u=>{calls.push(u);return response(u.split('=')[1]);}}); + assert.deepEqual(calls,['/index.html?v=b']); + const status=await(await c.match(core.completionKey(m))).json();assert.equal(status.offlineReady,true);assert.equal(status.failures,0); + await c.delete('/a.js?v=a');assert.equal((await core.status(m,c)).offlineReady,false); +}); + +test('completion retries with exponential backoff and keeps making progress past a failed asset',async()=>{ + const c=cache(),m=manifest(),delays=[];let tries=0; + const job=await core.complete(m,{cache:c,baseDelay:5,sleep:async ms=>delays.push(ms),fetchFn:async u=>{if(u.includes('a.js')&&++tries<=6)throw Error('drop');return response(u.split('=')[1]);}}); + assert.deepEqual(delays,[5,10]);assert.equal(job.offlineReady,true);assert.ok(await c.match('/index.html?v=b')); +}); + +test('overlapping layout and completion jobs share a six-download ceiling',async()=>{ + const c=cache(),m=manifest();m.groups.extra={background:true,files:[]}; + for(let i=0;i<15;i++){const p='/extra'+i+'.js';m.files[p]={h:String(i)};m.groups.extra.files.push(p);} + let active=0,peak=0;const calls=new Map(); + const fetchFn=async u=>{active++;peak=Math.max(active,peak);calls.set(u,(calls.get(u)||0)+1);await new Promise(r=>setTimeout(r,3));active--;return response(u.split('=')[1]);}; + await Promise.all([core.complete(m,{cache:c,fetchFn}),core.syncAssets(m,{cache:c,fetchFn,groups:['extra']})]); + assert.ok(peak<=6);assert.ok([...calls.values()].every(count=>count===1)); +}); diff --git a/frontend/asset-worker.test.js b/frontend/asset-worker.test.js index c3be060..7efcb6f 100644 --- a/frontend/asset-worker.test.js +++ b/frontend/asset-worker.test.js @@ -4,11 +4,11 @@ function environment(){ const storage={async keys(){return [...stores.keys()]},async delete(n){return stores.delete(n)},async open(n){if(!stores.has(n)){const map=new Map();stores.set(n,{async match(k){return map.get(typeof k==='string'?k:new URL(k.url).pathname+new URL(k.url).search)?.clone()},async put(k,r){map.set(k,r.clone())},async keys(){return [...map.keys()].map(k=>({url:'https://local'+k}))},async delete(k){return map.delete(typeof k==='string'?k:new URL(k.url).pathname+new URL(k.url).search)}})}return stores.get(n)}}; const sandbox={__BUILD_TAG__:'next',__ASSET_SYNC__:true,AssetSyncCore:core,importScripts:()=>{},caches:storage,URL,Response,Request,Headers,setTimeout,clearTimeout,crypto:require('node:crypto').webcrypto,console,fetch:async u=>{fetches.push(u);return u==='/api/manifest'?Response.json(m):new Response(u,{headers:{'X-Asset-Hash':m.files[u.split('?')[0]]?.h || ''}})},self:{location:{origin:'https://local'},addEventListener:(t,f)=>listeners[t]=f,skipWaiting:()=>skips.push(1),clients:{claim:async()=>{},matchAll:async()=>[{id:'old',postMessage:x=>messages.push(x)}]}}}; vm.runInNewContext(fs.readFileSync(__dirname+'/sw.js','utf8'),sandbox); - async function dispatch(t,e={}){let p;listeners[t]({...e,waitUntil:v=>p=v});await p} - async function request(url,mode='cors',clientId='new'){let p;listeners.fetch({request:{url:'https://local'+url,method:'GET',mode},clientId,respondWith:v=>p=v});return p} + async function dispatch(t,e={}){const tasks=[];listeners[t]({...e,waitUntil:v=>tasks.push(v)});await Promise.all(tasks)} + async function request(url,mode='cors',clientId='new'){let p;listeners.fetch({request:{url:'https://local'+url,method:'GET',mode},clientId,waitUntil:()=>{},respondWith:v=>p=v});return p} return {storage,m,skips,fetches,dispatch,request}; } -test('install does not activate or publish; CACHE_STATUS is honest; activation commits and exact requests self-heal',async()=>{const e=environment();await e.dispatch('install');assert.equal(e.skips.length,0);const cache=await e.storage.open(core.CACHE);assert.equal(await core.state(cache),null);let reply;await e.dispatch('message',{data:{type:'CACHE_STATUS'},ports:[{postMessage:r=>reply=r}]});assert.equal(reply.ready,true);await e.dispatch('activate');assert.equal((await core.state(cache)).current.buildTag,'next');assert.equal((await e.request('/playlist/x','navigate')).headers.get('X-Asset-Hash'),'index');await cache.delete('/app.js?v=app');await e.dispatch('message',{data:{type:'CACHE_STATUS'},ports:[{postMessage:r=>reply=r}]});assert.equal(reply.ready,false);assert.equal(reply.missing,1);assert.equal(reply.version,'next');assert.equal((await e.request('/app.js')).headers.get('X-Asset-Hash'),'app');assert.ok(await cache.match('/app.js?v=app'));assert.equal(await cache.match('/app.js?v=stale'),undefined);assert.equal((await e.request('/app.js?v=stale')).status,409);}); +test('install does not activate or publish; CACHE_STATUS is honest; activation commits and exact requests self-heal',async()=>{const e=environment();await e.dispatch('install');assert.equal(e.skips.length,0);const cache=await e.storage.open(core.CACHE);assert.equal(await core.state(cache),null);let reply;await e.dispatch('message',{data:{type:'CACHE_STATUS'},ports:[{postMessage:r=>reply=r}]});assert.equal(reply.ready,true);await e.dispatch('activate');assert.equal((await core.state(cache)).current.buildTag,'next');assert.equal((await e.request('/playlist/x','navigate')).headers.get('X-Asset-Hash'),'index');await cache.delete('/app.js?v=app');assert.equal((await core.status(e.m,cache)).ready,false);await e.dispatch('message',{data:{type:'CACHE_STATUS'},ports:[{postMessage:r=>reply=r}]});assert.equal((await core.status(e.m,cache)).offlineReady,true);assert.equal(reply.version,'next');assert.equal((await e.request('/app.js')).headers.get('X-Asset-Hash'),'app');assert.ok(await cache.match('/app.js?v=app'));assert.equal(await cache.match('/app.js?v=stale'),undefined);assert.equal((await e.request('/app.js?v=stale')).status,409);}); test('reported playback defers explicit activation until pause; install never calls skipWaiting',async()=>{const e=environment(),source={id:'p'};await e.dispatch('message',{data:{type:'SKIP_WAITING'}});assert.equal(e.skips.length,0);await e.dispatch('message',{source,data:{type:'PLAYING',value:true}});await e.dispatch('message',{source,data:{type:'SKIP_WAITING'}});assert.equal(e.skips.length,0);await e.dispatch('message',{source,data:{type:'PLAYING',value:false}});assert.equal(e.skips.length,1);}); test('staged same-contract responses use verified N-1 without poisoning the new URL',async()=>{ @@ -17,9 +17,10 @@ test('staged same-contract responses use verified N-1 without poisoning the new e.m.groups.core.contract=1; e.m.groups['feature:extra']={files:['/extra.js'],contract:1,background:true}; e.m.files['/extra.js']={h:'new'}; await cache.put(core.STATE,Response.json({current:old})); await cache.put('/extra.js?v=old',new Response('old body',{headers:{'X-Asset-Hash':'old'}})); await e.dispatch('install'); await e.dispatch('activate'); + assert.ok(await cache.match('/extra.js?v=new'));await cache.delete('/extra.js?v=new'); assert.equal(await cache.match('/extra.js?v=new'),undefined); const stale=await e.request('/extra.js?v=new'); assert.equal(await stale.text(),'old body');assert.equal(stale.headers.get('X-Asset-Hash'),'old');assert.equal(await cache.match('/extra.js?v=new'),undefined); - await e.dispatch('message',{data:{type:'WARM_ASSETS',saveData:true}}); assert.equal(await cache.match('/extra.js?v=new'),undefined); + await e.dispatch('message',{data:{type:'WARM_ASSETS',saveData:true}}); assert.ok(await cache.match('/extra.js?v=new')); await e.dispatch('message',{data:{type:'WARM_ASSETS',saveData:false}}); assert.equal((await cache.match('/extra.js?v=new')).headers.get('X-Asset-Hash'),'new'); assert.equal((await e.request('/extra.js?v=new')).headers.get('X-Asset-Hash'),'new'); }); diff --git a/frontend/sw.js b/frontend/sw.js index 0551c1b..64a8461 100644 --- a/frontend/sw.js +++ b/frontend/sw.js @@ -264,6 +264,8 @@ self.addEventListener('fetch', (e) => { return; } + if (ASSET_SYNC && request.mode === 'navigate') e.waitUntil(completeAssets()); + // Only intercept GET/HEAD — let POST (sync endpoint) go through unmodified if (request.method !== 'GET' && request.method !== 'HEAD') return; @@ -300,7 +302,32 @@ self.addEventListener('fetch', (e) => { // ---- Message: handle SKIP_WAITING from the client ---- const playingClients = new Map(), deferredClients = new Set(); -let activeLayout = 'classic', layoutReply, warmTask; +let activeLayout = 'classic', layoutReply; +let completionTask; +function completeAssets(force = false) { + if (!ASSET_SYNC) return Promise.resolve(); + completionTask ||= (async () => { + const cache = await caches.open(AssetSyncCore.CACHE); + const m = await candidate(cache) || (await AssetSyncCore.state(cache))?.current; + if (!m || m.buildTag !== VERSION) return; + const job = await AssetSyncCore.complete(m, { cache, fetchFn: fetch, force, + notify: async status => { + for (const client of await self.clients.matchAll({ type: 'window' })) + client.postMessage({ type: 'OFFLINE_STATUS', ...status }); + } + }); + if (!job.offlineReady) { + // Chromium can wake this job without a page; Safari resumes on the next + // navigation/online message. No timers outside an event lifetime. + try { await self.registration.sync?.register('ytp-offline-complete'); } catch {} + } + })().catch(error => console.warn('[asset-sync] completion paused', error.message)) + .finally(() => { completionTask = null; }); + return completionTask; +} +self.addEventListener('sync', e => { + if (e.tag === 'ytp-offline-complete') e.waitUntil(completeAssets()); +}); self.addEventListener('message', (e) => { if (ASSET_SYNC && e.data?.type === 'LAYOUT' && e.source && typeof e.data.value === 'string') { activeLayout = e.data.value; @@ -315,14 +342,8 @@ self.addEventListener('message', (e) => { } })()); } - if (ASSET_SYNC && e.data?.type === 'WARM_ASSETS' && !e.data.saveData) { - // Boot, activation and multiple tabs can report idle together. Share the - // verified pool so they never download the same missing group twice. - warmTask ||= (async()=>{ - const cache=await caches.open(AssetSyncCore.CACHE), m=(await AssetSyncCore.state(cache))?.current; - if(m)await AssetSyncCore.syncAssets(m,{cache,fetchFn:fetch,groups:Object.keys(m.groups).filter(name=>m.groups[name].background)}); - })().catch(error=>console.warn('[asset-sync] idle warm interrupted',error.message)).finally(()=>{warmTask=null;}); - e.waitUntil(warmTask); + if (ASSET_SYNC && ['WARM_ASSETS', 'COMPLETE_ASSETS'].includes(e.data?.type)) { + e.waitUntil(completeAssets(e.data.type === 'COMPLETE_ASSETS')); } if (e.data && e.data.type === 'PLAYING' && e.source) { playingClients.set(e.source.id, !!e.data.value); @@ -339,6 +360,7 @@ self.addEventListener('message', (e) => { // deletes its cache on failure, but a cache can also be evicted under // storage pressure, so the files are actually checked. if (e.data && e.data.type === 'CACHE_STATUS') { + if (ASSET_SYNC) e.waitUntil(completeAssets()); e.waitUntil((async () => { if (ASSET_SYNC) { let reply = { ready: false, missing: -1, version: VERSION }; @@ -347,7 +369,7 @@ self.addEventListener('message', (e) => { const pending = await candidate(cache), m = pending || state?.current; if(m) reply = await AssetSyncCore.status(m,cache,{previous:pending ? state?.current : state?.previous,activeLayout:e.data.activeLayout || (pending ? m.activeLayout : activeLayout) || 'classic'}); } catch {} - reply.type = 'CACHE_STATUS'; reply.assetSync = true; + reply.type = 'CACHE_STATUS'; reply.assetSync = true; reply.completing = !!completionTask && !reply.offlineReady; if(e.ports && e.ports[0]) e.ports[0].postMessage(reply); else if(e.source) e.source.postMessage(reply); return; } @@ -551,6 +573,7 @@ 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}); + await completeAssets(true); try { console.info('[asset-sync] storage', await self.navigator.storage.estimate()); } catch {} } async function assetFetch(request, clientId) {