import fs from 'node:fs/promises'; import path from 'node:path'; import { createHash, randomUUID } from 'node:crypto'; import webpush from 'web-push'; const idFor = value => createHash('sha256').update(value).digest('hex'); const DAY = 86400000; export function validateSubscription(value) { const endpoint = new URL(value?.endpoint); const host = endpoint.hostname; const allowed = host === 'fcm.googleapis.com' || host === 'web.push.apple.com' || host.endsWith('.push.apple.com') || host === 'updates.push.services.mozilla.com' || host.endsWith('.push.services.mozilla.com') || host.endsWith('.notify.windows.com'); if (endpoint.protocol !== 'https:' || endpoint.username || endpoint.password || endpoint.hash || (endpoint.port && endpoint.port !== '443') || !allowed || endpoint.href.length > 2048) throw new Error('Unsupported push service endpoint.'); for (const [key,size] of [['p256dh',65],['auth',16]]) { if (typeof value.keys?.[key] !== 'string' || !/^[A-Za-z0-9_-]+={0,2}$/.test(value.keys[key]) || Buffer.from(value.keys[key],'base64url').length !== size) throw new Error('Invalid push subscription keys.'); } if (Buffer.from(value.keys.p256dh,'base64url')[0] !== 4) throw new Error('Invalid push public key.'); return {endpoint:endpoint.href,keys:{p256dh:value.keys.p256dh,auth:value.keys.auth}}; } export class PushService { constructor({directory,subject,send=webpush.sendNotification,now=Date.now,log=console.warn}) { this.directory=directory;this.subject=subject;this.send=send;this.now=now;this.log=log;this.operations=Promise.resolve();this.flushing=false; } async init(articles) { await fs.mkdir(this.directory,{recursive:true,mode:0o700}); const keyFile=path.join(this.directory,'vapid.json'); try { this.keys=JSON.parse(await fs.readFile(keyFile,'utf8')); } catch(error) { if(error.code!=='ENOENT') throw error; this.keys=webpush.generateVAPIDKeys(); await fs.writeFile(keyFile,JSON.stringify(this.keys),{flag:'wx',mode:0o600}); } // Validate VAPID configuration now, before accepting subscriptions. webpush.generateRequestDetails({endpoint:'https://fcm.googleapis.com/fcm/send/config-check'},null,{vapidDetails:{subject:this.subject,...this.keys}}); try { this.state=JSON.parse(await fs.readFile(path.join(this.directory,'state.json'),'utf8')); } catch(error) { if(error.code!=='ENOENT') throw error; // Existing stories establish the baseline; never send a startup backlog. this.state={version:1,known:articles.map(a=>a.url),subscriptions:{},pending:[]}; await this.persist(); } if(this.state.version!==1 || !Array.isArray(this.state.known) || !this.state.subscriptions || !Array.isArray(this.state.pending)) throw new Error('Invalid push state; restore the data directory from backup.'); await this.sync(articles); return this; } transact(fn) { const operation=this.operations.then(fn); this.operations=operation.catch(()=>{}); return operation; } async persist() { const temporary=path.join(this.directory,`.state-${randomUUID()}.tmp`); try {await fs.writeFile(temporary,JSON.stringify(this.state),{mode:0o600});await fs.rename(temporary,path.join(this.directory,'state.json'));} finally {await fs.rm(temporary,{force:true});} } config() { return {enabled:true,publicKey:this.keys.publicKey}; } async subscribe(value) { const subscription=validateSubscription(value); return this.transact(async()=>{ const id=idFor(subscription.endpoint); if(!this.state.subscriptions[id] && Object.keys(this.state.subscriptions).length>=10000) throw new Error('Subscription capacity reached.'); this.state.subscriptions[id]={subscription}; await this.persist(); }); } async unsubscribe(value) { const subscription=validateSubscription(value); return this.transact(async()=>{ const id=idFor(subscription.endpoint); const existing=this.state.subscriptions[id]; if(existing && existing.subscription.keys.auth===subscription.keys.auth && existing.subscription.keys.p256dh===subscription.keys.p256dh) { delete this.state.subscriptions[id]; this.state.pending=this.state.pending.filter(job=>job.subscriber!==id); await this.persist(); } }); } async sync(articles) { return this.transact(async()=>{ const known=new Set(this.state.known); let changed=false; for(const article of articles) { if(known.has(article.url)) continue; changed=true;known.add(article.url); const payload={title:article.title.slice(0,120),body:article.description.slice(0,200),url:article.url,tag:`article-${idFor(article.url).slice(0,24)}`}; for(const subscriber of Object.keys(this.state.subscriptions)) this.state.pending.push({id:randomUUID(),subscriber,payload,attempts:0,retryAt:0,createdAt:this.now()}); } // Removed unpublished stories must not be delivered from an old queue. const published=new Set(articles.map(a=>a.url)); const pending=this.state.pending.filter(job=>published.has(job.payload.url)&&this.now()-job.createdAtthis.state.pending.filter(job=>job.retryAt<=this.now()).slice(0,20).map(job=>({...job,subscription:this.state.subscriptions[job.subscriber]?.subscription}))); // Small parallel batches keep subscribing/unsubscribing responsive. for(let offset=0;offset{ let status='sent'; if(job.subscription) { try { await this.send(job.subscription,JSON.stringify(job.payload),{vapidDetails:{subject:this.subject,...this.keys},TTL:86400,timeout:10000,urgency:'normal',topic:job.payload.tag}); } catch(error) { status=[404,410].includes(error.statusCode)?'expired':'retry'; if(status==='retry') this.log(`Push delivery delayed (${error.statusCode || 'network error'}).`); } } await this.transact(async()=>{ const current=this.state.pending.find(item=>item.id===job.id); if(!current) return; if(status==='expired') { delete this.state.subscriptions[job.subscriber]; this.state.pending=this.state.pending.filter(item=>item.subscriber!==job.subscriber); } else if(status==='retry' && current.attempts<4 && this.now()-current.createdAtitem.id!==job.id); } await this.persist(); }); })); } } finally { this.flushing=false; } } } export function createPushHandler(service,{siteUrl}={}) { const limits=new Map(); const json=(res,status,body)=>{res.writeHead(status,{'Content-Type':'application/json','Cache-Control':'no-store'});res.end(JSON.stringify(body));}; return async(req,res,pathname)=>{ if(!pathname.startsWith('/api/push/')) return false; if(pathname==='/api/push/config' && req.method==='GET') {json(res,200,service?service.config():{enabled:false});return true;} if(!['/api/push/subscribe','/api/push/unsubscribe'].includes(pathname)) {json(res,404,{error:'Unknown notification endpoint.'});return true;} if(req.method!=='POST') {json(res,405,{error:'Use POST.'});return true;} if(!service) {json(res,503,{error:'Notifications are not configured yet.'});return true;} const expected=siteUrl ? new URL(siteUrl).origin : `http://${req.headers.host}`; if(req.headers.origin!==expected || req.headers['sec-fetch-site']==='cross-site') {json(res,403,{error:'Origin not allowed.'});return true;} if(!req.headers['content-type']?.startsWith('application/json')) {json(res,415,{error:'Expected JSON.'});return true;} const now=Date.now(); for(const [key,value] of limits) if(value.until30 || limits.size>10000) {json(res,429,{error:'Too many requests. Try again in a minute.'});return true;} limits.set(ip,limit); try { let body=''; for await(const chunk of req) {body+=chunk;if(Buffer.byteLength(body)>8192){json(res,413,{error:'Subscription too large.'});return true;}} const subscription=JSON.parse(body); if(pathname.endsWith('/unsubscribe')) await service.unsubscribe(subscription); else await service.subscribe(subscription); json(res,200,{ok:true}); } catch {json(res,400,{error:'Unable to save this subscription. Please try again.'});} return true; }; }