From 5b25ccaa8693e4147a028edc978bc3735e9560d8 Mon Sep 17 00:00:00 2001 From: Jonathan Sykes Date: Wed, 7 Oct 2026 23:56:46 +0800 Subject: [PATCH] Publish asset state atomically and reject unavailable versions --- frontend/asset-sync-core.js | 131 +++++++++++++++++++++++++------ frontend/asset-sync-core.test.js | 3 + frontend/asset-worker.test.js | 4 +- frontend/sw.js | 11 ++- 4 files changed, 119 insertions(+), 30 deletions(-) diff --git a/frontend/asset-sync-core.js b/frontend/asset-sync-core.js index d3c9f37..e2b68fb 100644 --- a/frontend/asset-sync-core.js +++ b/frontend/asset-sync-core.js @@ -1,25 +1,112 @@ /* Exact, resumable asset synchronization shared by pages and workers. */ -(function(root){ +(function (root) { 'use strict'; - const CACHE='ytplayer-assets', STATE='/__ytp_asset_state'; - const url=(path,file)=>path+'?v='+file.h; - function paths(m){return [...new Set(Object.values(m.groups).flatMap(g=>g.files))].sort();} - function blocking(m){return [...new Set(Object.entries(m.groups).filter(([n,g])=>n==='core'||!g.background).flatMap(([,g])=>g.files))].sort();} - function plan(m,keys){const held=new Set(keys);return {missing:paths(m).map(p=>url(p,m.files[p])).filter(u=>!held.has(u)),blocking:blocking(m).map(p=>url(p,m.files[p]))};} - function retained(current,previous){const keep=new Set(paths(current).map(p=>url(p,current.files[p])));if(previous)for(const p of paths(previous))if(current.files[p]&¤t.files[p].h!==previous.files[p].h)keep.add(url(p,previous.files[p]));return keep;} - async function status(m,cache){let missing=0;for(const p of blocking(m)){const r=await cache.match(url(p,m.files[p]));if(!r||r.headers.get('X-Asset-Hash')!==m.files[p].h)missing++;}return {ready:missing===0,missing,version:m.buildTag};} - async function syncAssets(m,{cache,fetchFn,concurrency=6,attempts=3}){ - const missing=[];for(const p of paths(m)){const u=url(p,m.files[p]),r=await cache.match(u);if(!r||r.headers.get('X-Asset-Hash')!==m.files[p].h)missing.push(u);} - const count=missing.length;let failure; - await Promise.all(Array.from({length:Math.min(6,concurrency,Math.max(1,count))},async()=>{ - while(missing.length&&!failure){const u=missing.shift();let error; - for(let i=0;ictl.abort(),30000);let r;try{r=await fetchFn(u,{credentials:'same-origin',signal:ctl.signal});}finally{clearTimeout(timer);}if(!r.ok||r.headers.get('X-Asset-Hash')!==u.split('v=')[1])throw Error('Asset hash mismatch: '+u);await cache.put(u,r);error=null;break;}catch(e){error=e;}} - if(error)failure=error; - } - }));if(failure)throw failure;if(!(await status(m,cache)).ready)throw Error('Incomplete blocking assets');return {refreshed:count,caches:1}; + const CACHE = 'ytplayer-assets'; + const STATE = '/__ytp_asset_state'; + const url = (path, file) => path + '?v=' + file.h; + + function paths(manifest) { + return [...new Set(Object.values(manifest.groups).flatMap(group => group.files))].sort(); } - async function state(cache){const r=await cache.match(STATE);return r?await r.json():null;} - async function commit(m,cache){if(!(await status(m,cache)).ready)throw Error('Incomplete blocking assets');const old=await state(cache);const previous=old?.current?.buildTag===m.buildTag?old.previous:old?.current;await cache.put(STATE,new Response(JSON.stringify({current:m,previous}),{headers:{'Content-Type':'application/json'}}));const keep=retained(m,previous);for(const req of await cache.keys()){const u=new URL(req.url);const k=u.pathname+u.search;if(u.searchParams.has('v')&&!keep.has(k))await cache.delete(k);}return {current:m,previous};} - const api={CACHE,STATE,url,paths,blocking,plan,retained,status,syncAssets,state,commit}; - if(typeof module!=='undefined'&&module.exports)module.exports=api;else root.AssetSyncCore=api; -})(typeof globalThis!=='undefined'?globalThis:this); + + function blocking(manifest, previous, activeLayout) { + return [...new Set(Object.entries(manifest.groups) + .filter(([name, group]) => name === 'core' || !group.background || + name === 'layout:' + activeLayout || + (previous?.groups[name] && previous.groups[name].contract !== group.contract)) + .flatMap(([, group]) => group.files))].sort(); + } + + function plan(manifest, keys) { + const held = new Set(keys); + return { + missing: paths(manifest).map(path => url(path, manifest.files[path])).filter(key => !held.has(key)), + blocking: blocking(manifest).map(path => url(path, manifest.files[path])), + }; + } + + function retained(current, previous) { + const keep = new Set(paths(current).map(path => url(path, current.files[path]))); + if (previous) for (const path of paths(previous)) { + if (current.files[path]?.h !== previous.files[path].h) keep.add(url(path, previous.files[path])); + } + return keep; + } + + async function state(cache) { + const response = await cache.match(STATE); + return response ? response.json() : null; + } + + async function status(manifest, cache, { previous, activeLayout } = {}) { + const found = await Promise.all(blocking(manifest, previous, activeLayout).map(async path => { + const response = await cache.match(url(path, manifest.files[path])); + 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 }; + } + + async function syncAssets(manifest, { cache, fetchFn, concurrency = 6, attempts = 3, activeLayout }) { + const previous = (await state(cache))?.current; + const missing = []; + for (const path of blocking(manifest, previous, activeLayout)) { + const key = url(path, manifest.files[path]); + const response = await cache.match(key); + if (!response || response.headers.get('X-Asset-Hash') !== manifest.files[path].h) missing.push(key); + } + 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) { + const key = missing.shift(); + let error; + for (let attempt = 0; attempt < Math.min(3, attempts); attempt++) { + const controller = new AbortController(); + const timer = setTimeout(() => controller.abort(), 30000); + try { + const response = await fetchFn(key, { credentials: 'same-origin', signal: controller.signal }); + if (!response.ok || response.headers.get('X-Asset-Hash') !== key.split('v=')[1]) { + throw new Error('Asset hash mismatch: ' + key); + } + // Keep the timeout through body consumption, not just response headers. + await cache.put(key, response); + error = null; + break; + } catch (err) { error = err; } + finally { clearTimeout(timer); } + } + if (error) failure = error; + } + })); + if (failure) throw failure; + if (!(await status(manifest, cache, { previous, activeLayout })).ready) throw new Error('Incomplete blocking assets'); + return { refreshed: count, caches: 1 }; + } + + async function commit(manifest, cache, { previousClients, activeLayout } = {}) { + const old = await state(cache); + if (!(await status(manifest, cache, { previous: old?.current, activeLayout })).ready) { + throw new Error('Incomplete blocking assets'); + } + const sameBuild = old?.current?.buildTag === manifest.buildTag; + const previous = sameBuild ? old.previous : old?.current; + const next = { + current: manifest, previous, + previousClients: previousClients || (sameBuild ? old.previousClients : []), + }; + // One state write publishes the complete build and its old-tab affinity. + await cache.put(STATE, new Response(JSON.stringify(next), { headers: { 'Content-Type': 'application/json' } })); + const keep = retained(manifest, previous); + for (const request of await cache.keys()) { + const parsed = new URL(request.url); + const key = parsed.pathname + parsed.search; + if (parsed.searchParams.has('v') && !keep.has(key)) await cache.delete(key); + } + return next; + } + + const api = { CACHE, STATE, url, paths, blocking, plan, retained, status, 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 1b62137..a34a871 100644 --- a/frontend/asset-sync-core.test.js +++ b/frontend/asset-sync-core.test.js @@ -7,3 +7,6 @@ const response=h=>new Response(h,{headers:{'X-Asset-Hash':h}}); test('exact diff shares unchanged URLs; retains N-1 and prunes N-2',()=>{assert.deepEqual(core.plan(manifest(),['/a.js?v=a']).missing,['/index.html?v=b']);assert.deepEqual([...core.retained(manifest('c'),manifest())].sort(),['/a.js?v=a','/a.js?v=c','/index.html?v=b']);}); test('interrupted sync resumes verified files and refuses a mismatched hash',async()=>{const c=cache();let fail=true;let calls=[];const fetchFn=async u=>{calls.push(u);if(u.includes('index')&&fail)return response('wrong');return response(u.split('=')[1])};await assert.rejects(core.syncAssets(manifest(),{cache:c,fetchFn}),/hash/);assert.ok(await c.match('/a.js?v=a'));fail=false;calls=[];await core.syncAssets(manifest(),{cache:c,fetchFn});assert.deepEqual(calls,['/index.html?v=b']);await c.delete('/a.js?v=a');assert.equal((await core.status(manifest(),c)).missing,1);}); test('commit is atomic and pruning keeps only current and previous changed versions',async()=>{const c=cache();await c.put('/a.js?v=z',response('z'));await core.syncAssets(manifest(),{cache:c,fetchFn:async u=>response(u.split('=')[1])});await core.commit(manifest(),c);await core.syncAssets(manifest('c'),{cache:c,fetchFn:async u=>response(u.split('=')[1])});await core.commit(manifest('c'),c);assert.ok(await c.match('/a.js?v=a'));assert.equal(await c.match('/a.js?v=z'),undefined);await core.syncAssets(manifest('d'),{cache:c,fetchFn:async u=>response(u.split('=')[1])});await core.commit(manifest('d'),c);assert.equal(await c.match('/a.js?v=a'),undefined);}); +test('downloads are capped at six and each failure gets three attempts',async()=>{const c=cache(),files={},list=[];for(let i=0;i<13;i++){const p='/'+i+'.js';files[p]={h:String(i)};list.push(p)}const m={buildTag:'pool',files,groups:{core:{files:list}}};let active=0,max=0;const tries={};await core.syncAssets(m,{cache:c,concurrency:20,fetchFn:async u=>{active++;max=Math.max(max,active);await new Promise(r=>setTimeout(r,5));active--;tries[u]=(tries[u]||0)+1;if(tries[u]<3)throw Error('drop');return response(u.split('=')[1])}});assert.equal(max,6);assert.ok(Object.values(tries).every(n=>n===3));}); +test('background files do not block, but active layout and contract changes do',()=>{const m=manifest();m.groups.extra={background:true,contract:2,files:['/extra.js']};m.files['/extra.js']={h:'e'};assert.deepEqual(core.blocking(m),['/a.js','/index.html']);assert.ok(core.blocking(m,{groups:{extra:{contract:1}}}).includes('/extra.js'));m.groups['layout:classic']={background:true,files:['/classic.css']};m.files['/classic.css']={h:'c'};assert.ok(core.blocking(m,null,'classic').includes('/classic.css'));}); +test('removed files survive one previous build for open tabs',()=>{const previous=manifest(),current=manifest('c');delete current.files['/a.js'];current.groups.core.files=['/index.html'];assert.ok(core.retained(current,previous).has('/a.js?v=a'));}); diff --git a/frontend/asset-worker.test.js b/frontend/asset-worker.test.js index f7cc057..667c887 100644 --- a/frontend/asset-worker.test.js +++ b/frontend/asset-worker.test.js @@ -2,11 +2,11 @@ const {test}=require('node:test');const assert=require('node:assert/strict');con function environment(){ const stores=new Map(),listeners={},messages=[],skips=[];const m={buildTag:'next',files:{'/index.html':{h:'index'},'/app.js':{h:'app'}},groups:{core:{files:['/index.html','/app.js']}}}; 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,crypto:require('node:crypto').webcrypto,console,fetch:async u=>u==='/api/manifest'?Response.json(m):new Response(u,{headers:{'X-Asset-Hash':u.split('v=')[1]}}),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)}]}}}; + const sandbox={__BUILD_TAG__:'next',__ASSET_SYNC__:true,AssetSyncCore:core,importScripts:()=>{},caches:storage,URL,Response,Request,Headers,crypto:require('node:crypto').webcrypto,console,fetch:async u=>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} return {storage,m,skips,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((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);}); +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('reported playback defers explicit activation until pause; install never calls skipWaiting',async()=>{const e=environment(),source={id:'p'};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);}); diff --git a/frontend/sw.js b/frontend/sw.js index f41b7ff..43e7d6b 100644 --- a/frontend/sw.js +++ b/frontend/sw.js @@ -306,7 +306,7 @@ self.addEventListener('message', (e) => { e.waitUntil((async () => { if (ASSET_SYNC) { let reply = { ready: false, missing: -1, version: VERSION }; - try { const cache = await caches.open(AssetSyncCore.CACHE); const m = await candidate(cache); if(m) reply = await AssetSyncCore.status(m,cache); } catch {} + try { const cache = await caches.open(AssetSyncCore.CACHE); const m = await candidate(cache) || (await AssetSyncCore.state(cache))?.current; if(m) reply = await AssetSyncCore.status(m,cache); } catch {} reply.type = 'CACHE_STATUS'; reply.assetSync = true; if(e.ports && e.ports[0]) e.ports[0].postMessage(reply); else if(e.source) e.source.postMessage(reply); return; @@ -488,9 +488,8 @@ async function activateAssets() { const cache = await caches.open(AssetSyncCore.CACHE), m = await candidate(cache); if (!m) throw Error('Missing candidate manifest'); const names = await migrateLegacy(cache,m), old = await AssetSyncCore.state(cache); - const state = await AssetSyncCore.commit(m,cache); - state.previousClients = (await self.clients.matchAll({type:'window'})).map(c=>c.id); - await cache.put(AssetSyncCore.STATE,new Response(JSON.stringify(state))); + const previousClients = (await self.clients.matchAll({type:'window'})).map(c=>c.id); + await AssetSyncCore.commit(m,cache,{previousClients}); // Legacy deletion occurs strictly after verified commit. await Promise.all(names.map(n => caches.delete(n))); for (const req of await cache.keys()) if(new URL(req.url).pathname.startsWith('/__ytp_candidate/')) await cache.delete(req); @@ -509,7 +508,7 @@ async function assetFetch(request, clientId) { if(!nav && state.previousClients?.includes(clientId) && state.previous?.files[u.pathname]) m = state.previous; let path = nav ? '/index.html' : u.pathname; let key; - if(u.searchParams.has('v') && !nav) { + if(u.searchParams.has('v') && state.current.files[u.pathname]) { key = u.pathname + u.search; if(state.previous?.legacyURLs && u.searchParams.get('v') === state.previous.buildTag) key = state.previous.legacyURLs[path] || key; } else if(m.files[path]) key = AssetSyncCore.url(path,m.files[path]); @@ -517,7 +516,7 @@ async function assetFetch(request, clientId) { const cached = await cache.match(key); const h = new URL(key,self.location.origin).searchParams.get('v'); if(cached && cached.headers.get('X-Asset-Hash') === h) return cached; - try { const r = await fetch(key); if(r.ok && r.headers.get('X-Asset-Hash') === h) await cache.put(key,r.clone()); return r; } catch { return new Response('Offline',{status:503}); } + try { const r = await fetch(key); if(r.ok && r.headers.get('X-Asset-Hash') !== h) return new Response('Asset version unavailable',{status:409}); if(r.ok && request.method !== 'HEAD') await cache.put(key,r.clone()); return r; } catch { return new Response('Offline',{status:503}); } } try { return await fetch(request); } catch { return new Response('Offline',{status:503}); } }