fix(subsUnion): stop blocking DDP on subscription union recompute
Every subscription add/change/remove recomputed the full geo union
over all 7k+ subscriptions with a synchronous turf.union chain,
freezing the single Node event loop (and thus DDP/HTTP) for minutes.
The same recompute also runs at startup ("Subs union outdated"),
so every restart froze the site too.
calcUnionAsync yields to the event loop periodically during the
union chain, and subsUnion.js now fires recomputes without blocking
Meteor.startup or the observer callbacks, serializing overlapping
triggers instead of stacking them.
This commit is contained in:
parent
23e5f26239
commit
b4e5511cd0
4 changed files with 102 additions and 14 deletions
|
|
@ -6,7 +6,7 @@ import SiteSettings from '/imports/api/SiteSettings/SiteSettings';
|
|||
import Perlin from 'loms.perlin';
|
||||
import './leaflet-workaround';
|
||||
import L from 'leaflet';
|
||||
import calcUnion from '/imports/ui/components/Maps/SubsUnion/Unify';
|
||||
import calcUnionAsync from '/imports/startup/server/calcUnionAsync';
|
||||
import { isMailServerMaster } from '/imports/startup/server/email';
|
||||
|
||||
// sudo apt-get install libcairo2-dev libjpeg-dev libgif-dev
|
||||
|
|
@ -37,9 +37,8 @@ Meteor.startup(async () => {
|
|||
|
||||
const process = async (isPublic) => {
|
||||
const subscribers = await Subscriptions.find().fetchAsync();
|
||||
const result = calcUnion(L, subscribers, isPublic ? addNoisy : noNoisy, true);
|
||||
const union = result[0];
|
||||
const bounds = result[1];
|
||||
const union = await calcUnionAsync(subscribers, isPublic ? addNoisy : noNoisy);
|
||||
const bounds = union === null ? null : L.geoJSON(union).getBounds();
|
||||
|
||||
const publicl = isPublic ? 'public' : 'private';
|
||||
const Publicl = publicl.replace(/\b\w/g, l => l.toUpperCase());
|
||||
|
|
@ -83,6 +82,40 @@ Meteor.startup(async () => {
|
|||
}
|
||||
};
|
||||
|
||||
const recreate = async () => { await process(true); await process(false); };
|
||||
|
||||
let recomputeRunning = false;
|
||||
let recomputePending = false;
|
||||
|
||||
// Runs recreate() serially: a recompute already in flight is never overlapped by
|
||||
// another one, bursts of subscription changes just mark it pending and get a single
|
||||
// extra pass once the current one finishes.
|
||||
const runRecreate = async () => {
|
||||
recomputeRunning = true;
|
||||
try {
|
||||
await recreate();
|
||||
} catch (e) {
|
||||
console.error('subsUnion recompute failed', e);
|
||||
} finally {
|
||||
recomputeRunning = false;
|
||||
if (recomputePending) {
|
||||
recomputePending = false;
|
||||
runRecreate();
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
// Fire-and-forget on purpose: recreate() can take a while over thousands of
|
||||
// subscriptions, and callers (Meteor.startup, observeAsync callbacks) must not
|
||||
// block on it.
|
||||
const scheduleRecreate = () => {
|
||||
if (recomputeRunning) {
|
||||
recomputePending = true;
|
||||
return;
|
||||
}
|
||||
runRecreate();
|
||||
};
|
||||
|
||||
// At startup, we check if it's necessary to calc subscriptions union again
|
||||
const currentUnion = await SiteSettings.findOneAsync({ name: 'subs-public-union' });
|
||||
const lastSubs = await Subscriptions.findOneAsync({}, { sort: { updatedAt: -1 } });
|
||||
|
|
@ -93,30 +126,27 @@ Meteor.startup(async () => {
|
|||
const lastSubsUpdated = lastSubs.updatedAt;
|
||||
if (lastUnionUpdated > lastSubsUpdated || !countUnionSubs || countSubs !== countUnionSubs.value) {
|
||||
console.log('Subs union outdated');
|
||||
await process(true);
|
||||
await process(false);
|
||||
scheduleRecreate();
|
||||
} else {
|
||||
console.log('Subs union up-to-date');
|
||||
}
|
||||
}
|
||||
|
||||
const recreate = async () => { await process(true); await process(false); };
|
||||
|
||||
await Subscriptions.find({ createdAt: { $gt: new Date() } }).observeAsync({
|
||||
added: async function newSubAdded() { // doc) {
|
||||
added: function newSubAdded() { // doc) {
|
||||
if (debug) console.log('Subs added so recreate union');
|
||||
await recreate();
|
||||
scheduleRecreate();
|
||||
}
|
||||
});
|
||||
|
||||
await Subscriptions.find().observeAsync({
|
||||
changed: async function subsChanged() { // updatedDoc, oldDoc) {
|
||||
changed: function subsChanged() { // updatedDoc, oldDoc) {
|
||||
if (debug) console.log('Subs changed so recreate union');
|
||||
await recreate();
|
||||
scheduleRecreate();
|
||||
},
|
||||
removed: async function subsRemoved() { // oldDoc) {
|
||||
removed: function subsRemoved() { // oldDoc) {
|
||||
if (debug) console.log('Subs removed so recreate union');
|
||||
await recreate();
|
||||
scheduleRecreate();
|
||||
}
|
||||
});
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue