Files

160 lines
9.0 KiB
JavaScript

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.createdAt<DAY);
if(pending.length!==this.state.pending.length) changed=true;
this.state.pending=pending;
if(changed){this.state.known=[...known];await this.persist();}
});
}
async flush() {
if(this.flushing) return;
this.flushing=true;
try {
const jobs=await this.transact(()=>this.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<jobs.length;offset+=5) {
await Promise.all(jobs.slice(offset,offset+5).map(async job=>{
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.createdAt<DAY) {
current.attempts++;current.retryAt=this.now()+Math.min(3600000,30000*2**current.attempts);
} else {
if(status==='retry') this.log('Push delivery abandoned after repeated failures.');
this.state.pending=this.state.pending.filter(item=>item.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.until<now) limits.delete(key);
const ip=req.socket.remoteAddress || 'unknown';
const limit=limits.get(ip)||{count:0,until:now+60000};
if(++limit.count>30 || 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;
};
}