Skip to content

Commit 644ed79

Browse files
committed
CLEAN - review 1 performance
1 parent 6354b96 commit 644ed79

10 files changed

Lines changed: 369 additions & 99 deletions

File tree

dist/omt-router.cjs

Lines changed: 1 addition & 1 deletion
Large diffs are not rendered by default.

dist/omt-router.js

Lines changed: 1 addition & 1 deletion
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

src/engines/router.js

Lines changed: 84 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -102,14 +102,15 @@ function getPreparedGraph(graph, costField, normalizedPenalties, sourceId, targe
102102

103103
let _engineWorker = null;
104104
let _engineWorkerRequestId = 0;
105-
let _activeEngineJob = null;
105+
const _engineWorkerJobs = new Map();
106106
let _engineWorkerPreparedIdInWorker = null;
107107
let _engineWorkerPreparedIdCounter = 0;
108108
let _engineWorkerStatus = {
109109
state: ENGINE_WORKER_STATES.IDLE,
110110
running: false,
111111
engineId: null,
112112
requestId: null,
113+
requestCount: 0,
113114
startedAt: null,
114115
lastError: null,
115116
};
@@ -130,10 +131,11 @@ function emitEngineWorkerStatus(nextState) {
130131
}
131132
}
132133

133-
function rejectActiveEngineJob(error) {
134-
const job = _activeEngineJob;
135-
_activeEngineJob = null;
136-
if (job) job.reject(error);
134+
function rejectAllEngineJobs(error) {
135+
for (const job of _engineWorkerJobs.values()) {
136+
job.reject(error);
137+
}
138+
_engineWorkerJobs.clear();
137139
}
138140

139141
function terminateEngineWorker() {
@@ -163,51 +165,58 @@ function ensureEngineWorker() {
163165
const message = event.data ?? {};
164166

165167
if (message.type === 'status') {
166-
if (_activeEngineJob && message.requestId !== _activeEngineJob.requestId) return;
168+
const job = message.requestId != null ? _engineWorkerJobs.get(message.requestId) : null;
169+
const activeCount = _engineWorkerJobs.size;
167170

168171
if (message.state === ENGINE_WORKER_STATES.RUNNING) {
169172
emitEngineWorkerStatus({
170173
state: ENGINE_WORKER_STATES.RUNNING,
171174
running: true,
172-
engineId: normalizeEngineId(message.engineId, _activeEngineJob?.engineId ?? null),
173-
requestId: message.requestId ?? _activeEngineJob?.requestId ?? null,
174-
startedAt: _activeEngineJob?.startedAt ?? Date.now(),
175+
engineId: normalizeEngineId(message.engineId, job?.engineId ?? null),
176+
requestId: job?.requestId ?? null,
177+
requestCount: activeCount,
178+
startedAt: job?.startedAt ?? Date.now(),
175179
lastError: null,
176180
});
177181
} else if (message.state === ENGINE_WORKER_STATES.ERROR) {
178182
emitEngineWorkerStatus({
179183
state: ENGINE_WORKER_STATES.ERROR,
180-
running: false,
184+
running: activeCount > 0,
181185
engineId: null,
182186
requestId: null,
187+
requestCount: activeCount,
183188
startedAt: null,
184189
lastError: message.error ?? 'engine worker error',
185190
});
186-
} else if (message.state === ENGINE_WORKER_STATES.IDLE && !_activeEngineJob) {
191+
} else if (message.state === ENGINE_WORKER_STATES.IDLE) {
187192
emitEngineWorkerStatus({
188-
state: ENGINE_WORKER_STATES.IDLE,
189-
running: false,
190-
engineId: null,
191-
requestId: null,
192-
startedAt: null,
193+
state: activeCount > 0 ? ENGINE_WORKER_STATES.RUNNING : ENGINE_WORKER_STATES.IDLE,
194+
running: activeCount > 0,
195+
engineId: activeCount > 0 ? _engineWorkerStatus.engineId : null,
196+
requestId: activeCount > 0 ? _engineWorkerStatus.requestId : null,
197+
requestCount: activeCount,
198+
startedAt: activeCount > 0 ? _engineWorkerStatus.startedAt : null,
199+
lastError: null,
193200
});
194201
}
195202
return;
196203
}
197204

198205
if (message.type !== 'result') return;
199-
if (!_activeEngineJob || message.requestId !== _activeEngineJob.requestId) return;
206+
const job = _engineWorkerJobs.get(message.requestId);
207+
if (!job) return;
200208

201-
const { resolve, reject } = _activeEngineJob;
202-
_activeEngineJob = null;
209+
_engineWorkerJobs.delete(message.requestId);
210+
const { resolve, reject } = job;
203211

204212
if (message.ok) {
205213
emitEngineWorkerStatus({
206-
state: ENGINE_WORKER_STATES.IDLE,
207-
running: false,
208-
engineId: null,
209-
requestId: null,
210-
startedAt: null,
214+
state: _engineWorkerJobs.size > 0 ? ENGINE_WORKER_STATES.RUNNING : ENGINE_WORKER_STATES.IDLE,
215+
running: _engineWorkerJobs.size > 0,
216+
engineId: _engineWorkerJobs.size > 0 ? _engineWorkerStatus.engineId : null,
217+
requestId: _engineWorkerJobs.size > 0 ? _engineWorkerStatus.requestId : null,
218+
requestCount: _engineWorkerJobs.size,
219+
startedAt: _engineWorkerJobs.size > 0 ? _engineWorkerStatus.startedAt : null,
211220
lastError: null,
212221
});
213222
resolve(message.result);
@@ -234,12 +243,13 @@ function ensureEngineWorker() {
234243
running: false,
235244
engineId: null,
236245
requestId: null,
246+
requestCount: 0,
237247
startedAt: null,
238248
lastError: message,
239249
});
240250

241251
terminateEngineWorker();
242-
rejectActiveEngineJob(error);
252+
rejectAllEngineJobs(error);
243253
};
244254

245255
return _engineWorker;
@@ -252,6 +262,32 @@ function getEngineWorkerPreparedId(prepared) {
252262
return prepared._engineWorkerPreparedId;
253263
}
254264

265+
function ensureWorkerPreparedBackup(prepared) {
266+
if (!prepared._engineWorkerBackup) {
267+
prepared._engineWorkerBackup = {
268+
adjPtr: prepared.adjPtr.slice(),
269+
adjTo: prepared.adjTo.slice(),
270+
adjCost: prepared.adjCost.slice(),
271+
revAdjPtr: prepared.revAdjPtr.slice(),
272+
revAdjFrom: prepared.revAdjFrom.slice(),
273+
revAdjCost: prepared.revAdjCost.slice(),
274+
};
275+
}
276+
return prepared._engineWorkerBackup;
277+
}
278+
279+
function restorePreparedFromWorkerBackup(prepared) {
280+
const backup = prepared._engineWorkerBackup;
281+
if (!backup) return;
282+
283+
prepared.adjPtr = backup.adjPtr;
284+
prepared.adjTo = backup.adjTo;
285+
prepared.adjCost = backup.adjCost;
286+
prepared.revAdjPtr = backup.revAdjPtr;
287+
prepared.revAdjFrom = backup.revAdjFrom;
288+
prepared.revAdjCost = backup.revAdjCost;
289+
}
290+
255291
function serializePreparedForEngineWorker(prepared) {
256292
const coordsX = new Float32Array(prepared.N);
257293
const coordsY = new Float32Array(prepared.N);
@@ -261,14 +297,16 @@ function serializePreparedForEngineWorker(prepared) {
261297
coordsY[i] = coords?.[1] ?? 0;
262298
}
263299

300+
ensureWorkerPreparedBackup(prepared);
301+
264302
return {
265303
preparedId: getEngineWorkerPreparedId(prepared),
266-
adjPtr: prepared.adjPtr.slice(),
267-
adjTo: prepared.adjTo.slice(),
268-
adjCost: prepared.adjCost.slice(),
269-
revAdjPtr: prepared.revAdjPtr.slice(),
270-
revAdjFrom: prepared.revAdjFrom.slice(),
271-
revAdjCost: prepared.revAdjCost.slice(),
304+
adjPtr: prepared.adjPtr,
305+
adjTo: prepared.adjTo,
306+
adjCost: prepared.adjCost,
307+
revAdjPtr: prepared.revAdjPtr,
308+
revAdjFrom: prepared.revAdjFrom,
309+
revAdjCost: prepared.revAdjCost,
272310
N: prepared.N,
273311
E: prepared.E,
274312
coordsX,
@@ -330,36 +368,40 @@ async function runEngineInWorker(selectedEngine, startId, endId, prepared, {
330368
} = {}) {
331369
const worker = ensureEngineWorker();
332370
if (!worker) throw makeEngineError('engine worker is unavailable', 'engine_worker_unavailable');
333-
if (_activeEngineJob) throw makeEngineError('engine worker is busy', 'engine_worker_busy');
334371

335372
const requestId = ++_engineWorkerRequestId;
336373
const startedAt = Date.now();
337374

338375
return await new Promise((resolve, reject) => {
339-
_activeEngineJob = {
376+
_engineWorkerJobs.set(requestId, {
340377
requestId,
341378
resolve,
342379
reject,
343380
engineId: selectedEngine,
344381
startedAt,
345-
};
382+
});
346383

347384
emitEngineWorkerStatus({
348385
state: ENGINE_WORKER_STATES.RUNNING,
349386
running: true,
350387
engineId: selectedEngine,
351388
requestId,
389+
requestCount: _engineWorkerJobs.size,
352390
startedAt,
353391
lastError: null,
354392
});
355393

356394
const preparedId = getEngineWorkerPreparedId(prepared);
357395
if (_engineWorkerPreparedIdInWorker !== preparedId) {
358396
const serializedPrepared = serializePreparedForEngineWorker(prepared);
359-
worker.postMessage(
360-
{ type: 'prepare', prepared: serializedPrepared },
361-
getPreparedWorkerTransferables(serializedPrepared),
362-
);
397+
try {
398+
worker.postMessage(
399+
{ type: 'prepare', prepared: serializedPrepared },
400+
getPreparedWorkerTransferables(serializedPrepared),
401+
);
402+
} finally {
403+
restorePreparedFromWorkerBackup(prepared);
404+
}
363405
_engineWorkerPreparedIdInWorker = preparedId;
364406
}
365407

@@ -390,22 +432,24 @@ export function onEngineWorkerStatusChange(listener) {
390432
}
391433

392434
export function cancelRunningEngine(reason = 'cancelled') {
393-
if (!_activeEngineJob) return false;
435+
if (_engineWorkerJobs.size === 0) return false;
394436

395437
emitEngineWorkerStatus({
396438
state: ENGINE_WORKER_STATES.CANCELLING,
397439
running: true,
398440
lastError: reason,
441+
requestCount: _engineWorkerJobs.size,
399442
});
400443

401444
terminateEngineWorker();
402-
rejectActiveEngineJob(makeEngineError(`routing cancelled: ${reason}`, 'engine_cancelled'));
445+
rejectAllEngineJobs(makeEngineError(`routing cancelled: ${reason}`, 'engine_cancelled'));
403446

404447
emitEngineWorkerStatus({
405448
state: ENGINE_WORKER_STATES.IDLE,
406449
running: false,
407450
engineId: null,
408451
requestId: null,
452+
requestCount: 0,
409453
startedAt: null,
410454
lastError: reason,
411455
});

src/graphs/graphBuilder.js

Lines changed: 64 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -751,7 +751,8 @@ export function buildGraph(tiles, mode) {
751751
* 1. For each tile, check PowerCache for key `graph:<mode>:<z>/<x>/<y>`.
752752
* 2. On a cache miss, dispatch `{ op: 'parse-tile', url, x, y, z, mode }`
753753
* to a worker. The worker fetches the tile and calls parseTile itself.
754-
* 3. All tile parses are fanned out concurrently with Promise.all.
754+
* 3. Tile parses are dispatched in limited concurrent batches so large tile
755+
* sets are backpressured and fatal fetch failures can stop remaining work.
755756
* 4. Once every batch is resolved, run the sequential mergeSegments step on
756757
* the calling thread (node/edge deduplication cannot be parallelised).
757758
*
@@ -761,49 +762,75 @@ export function buildGraph(tiles, mode) {
761762
* @param {{
762763
* pool: import('@powerpool/powerpool').PowerPool,
763764
* cache: import('@powerpool/powercache').PowerCache,
764-
* ttl?: number
765+
* ttl?: number,
766+
* maxConcurrentTiles?: number
765767
* }} options
766-
* @returns {Promise<{ nodes: Map<number, object>, edges: Array<object>, nodeIndex: Map<string, number> }>}
768+
* @returns {Promise<{ nodes: Map<number, object>, edges: Array<object>, nodeIndex: Map<string, number>, hasMissingTiles: boolean, missingTileErrors: Array<object> }>}
767769
*/
768-
export async function buildGraphAsync(tiles, mode, { pool, cache, ttl = 300_000 } = {}) {
770+
function isFatalTileError(error) {
771+
if (!error || typeof error.code !== 'string') return false;
772+
return (
773+
error.code === 'MissingAllowOriginHeader' ||
774+
error.code === 'NetworkError' ||
775+
error.code.startsWith('HTTP_')
776+
);
777+
}
778+
779+
export async function buildGraphAsync(tiles, mode, { pool, cache, ttl = 300_000, maxConcurrentTiles } = {}) {
769780
if (!ways[mode]) {
770781
throw new Error(`Unknown transport mode "${mode}". Valid values: car, pedestrian, bicycle.`);
771782
}
772783

773784
const cacheVersion = 'v2';
774-
const results = await Promise.allSettled(
775-
tiles.map(({ url, x, y, z }) =>
776-
cache.getOrSetAsync(
777-
`graph:${cacheVersion}:${mode}:${z}/${x}/${y}`,
778-
async () => {
779-
const response = await pool.postMessage(
780-
{ op: 'parse-tile', url, x, y, z, mode },
781-
undefined,
782-
{ awaitResponse: true, timeout: 10_000 }
785+
const tileBatchSize = Math.max(
786+
1,
787+
typeof maxConcurrentTiles === 'number' && maxConcurrentTiles > 0
788+
? Math.min(maxConcurrentTiles, tiles.length)
789+
: Math.min(8, tiles.length, typeof pool?.maxSize === 'number' ? pool.maxSize : 4)
790+
);
791+
792+
const buildTileSegment = ({ url, x, y, z }) =>
793+
cache.getOrSetAsync(
794+
`graph:${cacheVersion}:${mode}:${z}/${x}/${y}`,
795+
async () => {
796+
const response = await pool.postMessage(
797+
{ op: 'parse-tile', url, x, y, z, mode },
798+
undefined,
799+
{ awaitResponse: true, timeout: 10_000 }
800+
);
801+
802+
if (response?.fetchFailed) {
803+
const reasonCode = response?.fetchError?.code;
804+
const reasonMessage = response?.fetchError?.message;
805+
const err = new Error(
806+
reasonMessage
807+
? `tile fetch failed for ${z}/${x}/${y}: ${reasonMessage}`
808+
: `tile fetch failed for ${z}/${x}/${y}`
783809
);
810+
err.code = reasonCode ?? 'TileFetchFailed';
811+
err.tile = { z, x, y, url };
812+
throw err;
813+
}
784814

785-
if (response?.fetchFailed) {
786-
const reasonCode = response?.fetchError?.code;
787-
const reasonMessage = response?.fetchError?.message;
788-
const err = new Error(
789-
reasonMessage
790-
? `tile fetch failed for ${z}/${x}/${y}: ${reasonMessage}`
791-
: `tile fetch failed for ${z}/${x}/${y}`
792-
);
793-
err.code = reasonCode ?? 'TileFetchFailed';
794-
err.tile = { z, x, y, url };
795-
throw err;
796-
}
815+
const payload = response?.output ?? response;
816+
return payload instanceof ArrayBuffer || ArrayBuffer.isView(payload)
817+
? u82o(payload)
818+
: payload;
819+
},
820+
{ ttl }
821+
);
797822

798-
const payload = response?.output ?? response;
799-
return payload instanceof ArrayBuffer || ArrayBuffer.isView(payload)
800-
? u82o(payload)
801-
: payload;
802-
},
803-
{ ttl }
804-
)
805-
)
806-
);
823+
const results = [];
824+
825+
for (let i = 0; i < tiles.length; i += tileBatchSize) {
826+
const batch = tiles.slice(i, i + tileBatchSize).map((tile) => buildTileSegment(tile));
827+
const settled = await Promise.allSettled(batch);
828+
results.push(...settled);
829+
830+
if (settled.some((result) => result.status === 'rejected' && isFatalTileError(result.reason))) {
831+
break;
832+
}
833+
}
807834

808835
const accumulator = createGraphAccumulator(mode);
809836
let hasMissingTiles = false;
@@ -826,6 +853,9 @@ export async function buildGraphAsync(tiles, mode, { pool, cache, ttl = 300_000
826853
const graph = finalizeGraph(accumulator);
827854
graph.hasMissingTiles = hasMissingTiles;
828855
graph.missingTileErrors = missingTileErrors;
856+
// `hasMissingTiles` flags that some tile fetch/parses failed. Callers can still
857+
// evaluate a route against the partial graph, but should avoid caching it as
858+
// a complete graph that may later produce false negatives.
829859
return graph;
830860
}
831861

0 commit comments

Comments
 (0)