diff --git a/.changeset/cold-apricots-poke.md b/.changeset/cold-apricots-poke.md new file mode 100644 index 0000000000..b050c8082b --- /dev/null +++ b/.changeset/cold-apricots-poke.md @@ -0,0 +1,5 @@ +--- +'@backstage/plugin-catalog-backend': patch +--- + +Changed the processing loop task pipeline implementation from recursive to iterative diff --git a/plugins/catalog-backend/src/processing/TaskPipeline.test.ts b/plugins/catalog-backend/src/processing/TaskPipeline.test.ts index f814e8a608..43ec118bdb 100644 --- a/plugins/catalog-backend/src/processing/TaskPipeline.test.ts +++ b/plugins/catalog-backend/src/processing/TaskPipeline.test.ts @@ -14,7 +14,7 @@ * limitations under the License. */ -import { startTaskPipeline } from './TaskPipeline'; +import { startTaskPipeline, createBarrier } from './TaskPipeline'; function createLimitedLoader(count: number, loadDelay?: number) { const items = new Array(count).fill(0).map((_, index) => index); @@ -119,3 +119,88 @@ describe('startTaskPipeline', () => { }).toThrow('must be lower'); }); }); + +describe('createBarrier', () => { + const tick = (millis: number) => + new Promise(resolve => setTimeout(resolve, millis)); + + it('abandons a wait after the timeout expires', async () => { + const abortController = new AbortController(); + const signal = abortController.signal; + const barrier = createBarrier({ waitTimeoutMillis: 100, signal }); + + const fn1 = jest.fn(); + barrier.wait().then(fn1); + + await tick(0); + expect(fn1).not.toHaveBeenCalled(); + + await tick(50); + expect(fn1).not.toHaveBeenCalled(); + + // start a new wait mid-way through the timeout + // should NOT resolve when the first one times out + const fn2 = jest.fn(); + barrier.wait().then(fn2); + + await tick(0); + expect(fn2).not.toHaveBeenCalled(); + + await tick(50); + expect(fn1).toHaveBeenCalledTimes(1); + expect(fn2).not.toHaveBeenCalled(); + + await tick(50); + expect(fn1).toHaveBeenCalledTimes(1); + expect(fn2).toHaveBeenCalledTimes(1); + }); + + it('abandons a wait after aborted', async () => { + const abortController = new AbortController(); + const signal = abortController.signal; + const barrier = createBarrier({ waitTimeoutMillis: 100, signal }); + + const fn1 = jest.fn(); + barrier.wait().then(fn1); + + // should resolve immediately, not after timeout + await tick(0); + expect(fn1).not.toHaveBeenCalled(); + abortController.abort(); + await tick(0); + expect(fn1).toHaveBeenCalledTimes(1); + + // subsequent waits should be immediate no matter what + const fn2 = jest.fn(); + barrier.wait().then(fn2); + await tick(0); + expect(fn2).toHaveBeenCalledTimes(1); + }); + + it('release immediately unblocks all waits', async () => { + const abortController = new AbortController(); + const signal = abortController.signal; + const barrier = createBarrier({ waitTimeoutMillis: 100, signal }); + + const fn1 = jest.fn(); + barrier.wait().then(fn1); + + await tick(50); + expect(fn1).not.toHaveBeenCalled(); + + // start a new wait mid-way through the timeout + // SHOULD resolve when releasing + const fn2 = jest.fn(); + barrier.wait().then(fn2); + + await tick(0); + expect(fn1).not.toHaveBeenCalled(); + expect(fn2).not.toHaveBeenCalled(); + + barrier.release(); + + await tick(0); + expect(fn1).toHaveBeenCalledTimes(1); + expect(fn2).toHaveBeenCalledTimes(1); + }); +}); diff --git a/plugins/catalog-backend/src/processing/TaskPipeline.ts b/plugins/catalog-backend/src/processing/TaskPipeline.ts index 4afe2408a8..44030be9e0 100644 --- a/plugins/catalog-backend/src/processing/TaskPipeline.ts +++ b/plugins/catalog-backend/src/processing/TaskPipeline.ts @@ -72,50 +72,112 @@ export function startTaskPipeline(options: Options) { throw new Error('lowWatermark must be lower than highWatermark'); } - let loading = false; - let stopped = false; - let inFlightCount = 0; + // State is in an object so that it can be stably referenced from within + // callbacks below + const state = { inFlightCount: 0 }; + const abortController = new AbortController(); + const abortSignal = abortController.signal; - async function maybeLoadMore() { - if (stopped || loading || inFlightCount > lowWatermark) { - return; - } + const barrier = createBarrier({ + waitTimeoutMillis: pollingIntervalMs, + signal: abortSignal, + }); - // Once we hit the low watermark we load in enough items to reach the high watermark - loading = true; - const loadCount = highWatermark - inFlightCount; - const loadedItems = await loadTasks(loadCount); - loading = false; - - // We might not reach the high watermark here, in case there weren't enough items to load - inFlightCount += loadedItems.length; - loadedItems.forEach(item => { - processTask(item).finally(() => { - if (stopped) { - return; + async function pipelineLoop() { + while (!abortSignal.aborted) { + if (state.inFlightCount <= lowWatermark) { + const loadCount = highWatermark - state.inFlightCount; + const loadedItems = await Promise.resolve() + .then(() => loadTasks(loadCount)) + .catch(() => { + // Silently swallow errors and go back to sleep to try again; we + // delegate to the loadTasks function itself to catch errors and log + // if it so desires + return []; + }); + if (loadedItems.length && !abortSignal.aborted) { + state.inFlightCount += loadedItems.length; + for (const item of loadedItems) { + Promise.resolve() + .then(() => processTask(item)) + .catch(() => { + // Silently swallow errors and go back to sleep to try again; we + // delegate to the processTask function itself to catch errors + // and log if it so desires + }) + .finally(() => { + state.inFlightCount -= 1; + barrier.release(); + }); + } } + } - // For each item we complete we check if it's time to load more - inFlightCount -= 1; - maybeLoadMore(); - }); - }); - - // We might have processed some tasks while we where loading, so check if we can load more - if (loadedItems.length > 1) { - maybeLoadMore(); + await barrier.wait(); } } - // This interval makes sure that we load in new items if the loop runs - // dry because of the lack of available tasks. As long as there are - // enough items to process this will be a noop. - const intervalId = setInterval(() => { - maybeLoadMore(); - }, pollingIntervalMs); + pipelineLoop().catch(error => { + // This should be impossible, but if it did happen, it would signal a + // programming error inside the loop (errors should definitely be caught + // inside of it). Let's rethrow with more information, and let it be caught + // by the process' uncaught exception handler, which will log the occurrence + // at a high level. + throw new Error(`Unexpected error in processing pipeline loop`, error); + }); return () => { - stopped = true; - clearInterval(intervalId); + abortController.abort(); + barrier.destroy(); + }; +} + +/** + * Creates a barrier with a timeout, that can be awaited or prematurely + * released either manually or by an abort signal. + */ +export function createBarrier(options: { + waitTimeoutMillis: number; + signal: AbortSignal; +}): { + wait: () => Promise; + release: () => void; + destroy: () => void; +} { + const { waitTimeoutMillis, signal } = options; + const resolvers = new Set<() => void>(); + + function wait() { + if (signal.aborted || !(waitTimeoutMillis > 0)) { + return Promise.resolve(); + } + + return new Promise(resolve => { + const timeoutHandle = setTimeout(done, waitTimeoutMillis); + + function done() { + resolvers.delete(done); + clearTimeout(timeoutHandle); + resolve(); + } + + resolvers.add(done); + }); + } + + function release() { + const resolversToCall = new Set(resolvers); + resolvers.clear(); + for (const resolver of resolversToCall) { + resolver(); + } + } + + signal.addEventListener('abort', release); + + return { + wait, + release, + destroy: () => signal.removeEventListener('abort', release), }; }