Publish asset state atomically and reject unavailable versions

This commit is contained in:
Jonathan Sykes
2026-10-07 23:56:46 +08:00
parent 54a07a0a2a
commit 5b25ccaa86
4 changed files with 119 additions and 30 deletions

View File

@@ -1,25 +1,112 @@
/* Exact, resumable asset synchronization shared by pages and workers. */ /* Exact, resumable asset synchronization shared by pages and workers. */
(function(root){ (function (root) {
'use strict'; 'use strict';
const CACHE='ytplayer-assets', STATE='/__ytp_asset_state'; const CACHE = 'ytplayer-assets';
const url=(path,file)=>path+'?v='+file.h; const STATE = '/__ytp_asset_state';
function paths(m){return [...new Set(Object.values(m.groups).flatMap(g=>g.files))].sort();} const url = (path, file) => path + '?v=' + file.h;
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 paths(manifest) {
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]&&current.files[p].h!==previous.files[p].h)keep.add(url(p,previous.files[p]));return keep;} return [...new Set(Object.values(manifest.groups).flatMap(group => group.files))].sort();
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;i<attempts;i++){try{const ctl=new AbortController();const timer=setTimeout(()=>ctl.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};
} }
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};} function blocking(manifest, previous, activeLayout) {
const api={CACHE,STATE,url,paths,blocking,plan,retained,status,syncAssets,state,commit}; return [...new Set(Object.entries(manifest.groups)
if(typeof module!=='undefined'&&module.exports)module.exports=api;else root.AssetSyncCore=api; .filter(([name, group]) => name === 'core' || !group.background ||
})(typeof globalThis!=='undefined'?globalThis:this); 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);

View File

@@ -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('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('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('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'));});

View File

@@ -2,11 +2,11 @@ const {test}=require('node:test');const assert=require('node:assert/strict');con
function environment(){ 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 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 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); 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 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 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}; 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);}); 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);});

View File

@@ -306,7 +306,7 @@ self.addEventListener('message', (e) => {
e.waitUntil((async () => { e.waitUntil((async () => {
if (ASSET_SYNC) { if (ASSET_SYNC) {
let reply = { ready: false, missing: -1, version: VERSION }; 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; 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); if(e.ports && e.ports[0]) e.ports[0].postMessage(reply); else if(e.source) e.source.postMessage(reply);
return; return;
@@ -488,9 +488,8 @@ async function activateAssets() {
const cache = await caches.open(AssetSyncCore.CACHE), m = await candidate(cache); const cache = await caches.open(AssetSyncCore.CACHE), m = await candidate(cache);
if (!m) throw Error('Missing candidate manifest'); if (!m) throw Error('Missing candidate manifest');
const names = await migrateLegacy(cache,m), old = await AssetSyncCore.state(cache); const names = await migrateLegacy(cache,m), old = await AssetSyncCore.state(cache);
const state = await AssetSyncCore.commit(m,cache); const previousClients = (await self.clients.matchAll({type:'window'})).map(c=>c.id);
state.previousClients = (await self.clients.matchAll({type:'window'})).map(c=>c.id); await AssetSyncCore.commit(m,cache,{previousClients});
await cache.put(AssetSyncCore.STATE,new Response(JSON.stringify(state)));
// Legacy deletion occurs strictly after verified commit. // Legacy deletion occurs strictly after verified commit.
await Promise.all(names.map(n => caches.delete(n))); 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); 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; if(!nav && state.previousClients?.includes(clientId) && state.previous?.files[u.pathname]) m = state.previous;
let path = nav ? '/index.html' : u.pathname; let path = nav ? '/index.html' : u.pathname;
let key; let key;
if(u.searchParams.has('v') && !nav) { if(u.searchParams.has('v') && state.current.files[u.pathname]) {
key = u.pathname + u.search; key = u.pathname + u.search;
if(state.previous?.legacyURLs && u.searchParams.get('v') === state.previous.buildTag) key = state.previous.legacyURLs[path] || key; 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]); } 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 cached = await cache.match(key);
const h = new URL(key,self.location.origin).searchParams.get('v'); const h = new URL(key,self.location.origin).searchParams.get('v');
if(cached && cached.headers.get('X-Asset-Hash') === h) return cached; 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}); } try { return await fetch(request); } catch { return new Response('Offline',{status:503}); }
} }