Complete verified offline assets in a resumable worker job
This commit is contained in:
@@ -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);
|
||||
|
||||
@@ -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));
|
||||
});
|
||||
|
||||
@@ -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');
|
||||
});
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user