From b1254b40a8b78782acc4ecace2e118b8362f3b75 Mon Sep 17 00:00:00 2001 From: Eric Peterson Date: Tue, 1 Oct 2024 12:57:15 +0200 Subject: [PATCH 1/9] Add root lifecycle shutdown hook to clean up DB connections Signed-off-by: Eric Peterson --- .../entrypoints/database/DatabaseManager.ts | 28 +++++++++++++++++++ .../database/databaseServiceFactory.ts | 12 ++++++-- 2 files changed, 38 insertions(+), 2 deletions(-) diff --git a/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts b/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts index 46d4ee8cf7..6ed5545d35 100644 --- a/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts +++ b/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts @@ -89,6 +89,27 @@ export class DatabaseManagerImpl { return { getClient, migrations: { skip } }; } + /** + * Method to be called during shutdown to destroy all known connections. + */ + async shutdown(deps: { logger: LoggerService }): Promise { + const pluginIds = Array.from(this.databaseCache.keys()); + await Promise.all( + pluginIds.map(async pluginId => { + const connection = await this.databaseCache.get(pluginId); + if (connection) { + await connection.destroy().catch((error: unknown) => { + deps.logger.error( + `Problem closing database connection for ${pluginId}: ${stringifyError( + error, + )}`, + ); + }); + } + }), + ); + } + /** * Provides the client type which should be used for a given plugin. * @@ -236,4 +257,11 @@ export class DatabaseManager { ): PluginDatabaseManager { return this.impl.forPlugin(pluginId, deps); } + + /** + * Method to be called during shutdown to destroy all known connections. + */ + async shutdown(deps: { logger: LoggerService }): Promise { + return this.impl.shutdown(deps); + } } diff --git a/packages/backend-defaults/src/entrypoints/database/databaseServiceFactory.ts b/packages/backend-defaults/src/entrypoints/database/databaseServiceFactory.ts index 167779e807..5e23e6ac4c 100644 --- a/packages/backend-defaults/src/entrypoints/database/databaseServiceFactory.ts +++ b/packages/backend-defaults/src/entrypoints/database/databaseServiceFactory.ts @@ -37,9 +37,11 @@ export const databaseServiceFactory = createServiceFactory({ lifecycle: coreServices.lifecycle, logger: coreServices.logger, pluginMetadata: coreServices.pluginMetadata, + rootLifecycle: coreServices.rootLifecycle, + rootLogger: coreServices.rootLogger, }, - async createRootContext({ config }) { - return config.getOptional('backend.database') + async createRootContext({ config, rootLifecycle, rootLogger }) { + const databaseManager = config.getOptional('backend.database') ? DatabaseManager.fromConfig(config) : DatabaseManager.fromConfig( new ConfigReader({ @@ -48,6 +50,12 @@ export const databaseServiceFactory = createServiceFactory({ }, }), ); + + rootLifecycle.addShutdownHook(async () => { + await databaseManager.shutdown({ logger: rootLogger }); + }); + + return databaseManager; }, async factory({ pluginMetadata, lifecycle, logger }, databaseManager) { return databaseManager.forPlugin(pluginMetadata.getId(), { From 8bfe9f0f61aad548b6d129610b05b56aa9f93f2b Mon Sep 17 00:00:00 2001 From: Eric Peterson Date: Tue, 1 Oct 2024 12:57:43 +0200 Subject: [PATCH 2/9] Ensure root lifecycle shutdown hooks happen after plugins' Signed-off-by: Eric Peterson --- .../src/wiring/BackendInitializer.ts | 24 ++++++++- .../lifecycle/lifecycleServiceFactory.ts | 51 +++++++++++++++---- 2 files changed, 63 insertions(+), 12 deletions(-) diff --git a/packages/backend-app-api/src/wiring/BackendInitializer.ts b/packages/backend-app-api/src/wiring/BackendInitializer.ts index 86771fbb74..e580109987 100644 --- a/packages/backend-app-api/src/wiring/BackendInitializer.ts +++ b/packages/backend-app-api/src/wiring/BackendInitializer.ts @@ -407,6 +407,26 @@ export class BackendInitializer { // The startup failed, but we may still want to do cleanup so we continue silently } + // Get all plugins. + const pluginMap = new Map(); + for (const feature of this.#registrations) { + for (const r of feature.getRegistrations()) { + if (r.type === 'plugin') { + pluginMap.set(r.pluginId, true); + } + } + } + const allPluginIds = Array.from(pluginMap.keys()); + + // Iterate through all plugins and run their shutdown hooks. + await Promise.allSettled( + allPluginIds.map(async pluginId => { + const lifecycleService = await this.#getPluginLifecycleImpl(pluginId); + await lifecycleService.shutdown(); + }), + ); + + // Once all plugin shutdown hooks are done, run root shutdown hooks. const lifecycleService = await this.#getRootLifecycleImpl(); await lifecycleService.shutdown(); } @@ -437,7 +457,9 @@ export class BackendInitializer { async #getPluginLifecycleImpl( pluginId: string, - ): Promise }> { + ): Promise< + LifecycleService & { startup(): Promise; shutdown(): Promise } + > { const lifecycleService = await this.#serviceRegistry.get( coreServices.lifecycle, pluginId, diff --git a/packages/backend-defaults/src/entrypoints/lifecycle/lifecycleServiceFactory.ts b/packages/backend-defaults/src/entrypoints/lifecycle/lifecycleServiceFactory.ts index 192886ca61..02195d0cd0 100644 --- a/packages/backend-defaults/src/entrypoints/lifecycle/lifecycleServiceFactory.ts +++ b/packages/backend-defaults/src/entrypoints/lifecycle/lifecycleServiceFactory.ts @@ -22,7 +22,6 @@ import { LifecycleServiceStartupOptions, LoggerService, PluginMetadataService, - RootLifecycleService, coreServices, createServiceFactory, } from '@backstage/backend-plugin-api'; @@ -31,15 +30,19 @@ import { export class BackendPluginLifecycleImpl implements LifecycleService { constructor( private readonly logger: LoggerService, - private readonly rootLifecycle: RootLifecycleService, private readonly pluginMetadata: PluginMetadataService, ) {} #hasStarted = false; + #hasShutdown = false; #startupTasks: Array<{ hook: LifecycleServiceStartupHook; options?: LifecycleServiceStartupOptions; }> = []; + #shutdownTasks: Array<{ + hook: LifecycleServiceShutdownHook; + options?: LifecycleServiceShutdownOptions; + }> = []; addStartupHook( hook: LifecycleServiceStartupHook, @@ -77,11 +80,42 @@ export class BackendPluginLifecycleImpl implements LifecycleService { hook: LifecycleServiceShutdownHook, options?: LifecycleServiceShutdownOptions, ): void { + if (this.#hasShutdown) { + throw new Error('Attempted to add shutdown hook after shutdown'); + } const plugin = this.pluginMetadata.getId(); - this.rootLifecycle.addShutdownHook(hook, { - logger: options?.logger?.child({ plugin }) ?? this.logger, + const logger = options?.logger?.child({ plugin }) ?? this.logger; + this.#shutdownTasks.push({ + hook, + options: { + ...options, + logger, + }, }); } + + async shutdown(): Promise { + if (this.#hasShutdown) { + return; + } + this.#hasShutdown = true; + + this.logger.debug( + `Running ${this.#shutdownTasks.length} plugin shutdown tasks...`, + ); + + await Promise.all( + this.#shutdownTasks.map(async ({ hook, options }) => { + const logger = options?.logger ?? this.logger; + try { + await hook(); + logger.debug(`Plugin shutdown hook succeeded`); + } catch (error) { + logger.error(`Plugin shutdown hook failed, ${error}`); + } + }), + ); + } } /** @@ -97,14 +131,9 @@ export const lifecycleServiceFactory = createServiceFactory({ service: coreServices.lifecycle, deps: { logger: coreServices.logger, - rootLifecycle: coreServices.rootLifecycle, pluginMetadata: coreServices.pluginMetadata, }, - async factory({ rootLifecycle, logger, pluginMetadata }) { - return new BackendPluginLifecycleImpl( - logger, - rootLifecycle, - pluginMetadata, - ); + async factory({ logger, pluginMetadata }) { + return new BackendPluginLifecycleImpl(logger, pluginMetadata); }, }); From 1aaaf7033705b1199be90f3b3c9586a412171e49 Mon Sep 17 00:00:00 2001 From: Eric Peterson Date: Tue, 1 Oct 2024 14:10:16 +0200 Subject: [PATCH 3/9] Minor cleanup based on feedback. Signed-off-by: Eric Peterson --- .../src/wiring/BackendInitializer.ts | 13 ++++++++----- .../src/entrypoints/database/DatabaseManager.ts | 2 +- .../lifecycle/lifecycleServiceFactory.ts | 2 +- 3 files changed, 10 insertions(+), 7 deletions(-) diff --git a/packages/backend-app-api/src/wiring/BackendInitializer.ts b/packages/backend-app-api/src/wiring/BackendInitializer.ts index e580109987..3ea1200898 100644 --- a/packages/backend-app-api/src/wiring/BackendInitializer.ts +++ b/packages/backend-app-api/src/wiring/BackendInitializer.ts @@ -408,19 +408,18 @@ export class BackendInitializer { } // Get all plugins. - const pluginMap = new Map(); + const allPlugins = new Set(); for (const feature of this.#registrations) { for (const r of feature.getRegistrations()) { if (r.type === 'plugin') { - pluginMap.set(r.pluginId, true); + allPlugins.add(r.pluginId); } } } - const allPluginIds = Array.from(pluginMap.keys()); // Iterate through all plugins and run their shutdown hooks. await Promise.allSettled( - allPluginIds.map(async pluginId => { + [...allPlugins].map(async pluginId => { const lifecycleService = await this.#getPluginLifecycleImpl(pluginId); await lifecycleService.shutdown(); }), @@ -466,7 +465,11 @@ export class BackendInitializer { ); const service = lifecycleService as any; - if (service && typeof service.startup === 'function') { + if ( + service && + typeof service.startup === 'function' && + typeof service.shutdown === 'function' + ) { return service; } diff --git a/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts b/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts index 6ed5545d35..ff95c0562a 100644 --- a/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts +++ b/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts @@ -94,7 +94,7 @@ export class DatabaseManagerImpl { */ async shutdown(deps: { logger: LoggerService }): Promise { const pluginIds = Array.from(this.databaseCache.keys()); - await Promise.all( + await Promise.allSettled( pluginIds.map(async pluginId => { const connection = await this.databaseCache.get(pluginId); if (connection) { diff --git a/packages/backend-defaults/src/entrypoints/lifecycle/lifecycleServiceFactory.ts b/packages/backend-defaults/src/entrypoints/lifecycle/lifecycleServiceFactory.ts index 02195d0cd0..0ba812f5b4 100644 --- a/packages/backend-defaults/src/entrypoints/lifecycle/lifecycleServiceFactory.ts +++ b/packages/backend-defaults/src/entrypoints/lifecycle/lifecycleServiceFactory.ts @@ -111,7 +111,7 @@ export class BackendPluginLifecycleImpl implements LifecycleService { await hook(); logger.debug(`Plugin shutdown hook succeeded`); } catch (error) { - logger.error(`Plugin shutdown hook failed, ${error}`); + logger.error('Plugin shutdown hook failed', error); } }), ); From 8fa3ec895fd6d6d6c8bb5eca73b9f4a87791caaf Mon Sep 17 00:00:00 2001 From: Eric Peterson Date: Tue, 1 Oct 2024 14:23:22 +0200 Subject: [PATCH 4/9] Refactor DB Manager shutdown to be private Signed-off-by: Eric Peterson --- .../entrypoints/database/DatabaseManager.ts | 27 +++++++++++-------- .../database/databaseServiceFactory.ts | 11 +++----- 2 files changed, 19 insertions(+), 19 deletions(-) diff --git a/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts b/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts index ff95c0562a..9441c8bca9 100644 --- a/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts +++ b/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts @@ -19,6 +19,8 @@ import { LifecycleService, LoggerService, RootConfigService, + RootLifecycleService, + RootLoggerService, } from '@backstage/backend-plugin-api'; import { Config } from '@backstage/config'; import { stringifyError } from '@backstage/errors'; @@ -42,6 +44,8 @@ function pluginPath(pluginId: string): string { */ export type DatabaseManagerOptions = { migrations?: DatabaseService['migrations']; + rootLogger?: RootLoggerService; + rootLifecycle?: RootLifecycleService; }; /** @@ -53,7 +57,15 @@ export class DatabaseManagerImpl { private readonly connectors: Record, private readonly options?: DatabaseManagerOptions, private readonly databaseCache: Map> = new Map(), - ) {} + ) { + // If a rootLifecycle service was provided, register a shutdown hook to + // clean up any database connections. + if (options?.rootLifecycle !== undefined) { + options.rootLifecycle.addShutdownHook(async () => { + await this.shutdown({ logger: options.rootLogger }); + }); + } + } /** * Generates a PluginDatabaseManager for consumption by plugins. @@ -90,16 +102,16 @@ export class DatabaseManagerImpl { } /** - * Method to be called during shutdown to destroy all known connections. + * Destroys all known connections. */ - async shutdown(deps: { logger: LoggerService }): Promise { + private async shutdown(deps?: { logger?: LoggerService }): Promise { const pluginIds = Array.from(this.databaseCache.keys()); await Promise.allSettled( pluginIds.map(async pluginId => { const connection = await this.databaseCache.get(pluginId); if (connection) { await connection.destroy().catch((error: unknown) => { - deps.logger.error( + deps?.logger?.error( `Problem closing database connection for ${pluginId}: ${stringifyError( error, )}`, @@ -257,11 +269,4 @@ export class DatabaseManager { ): PluginDatabaseManager { return this.impl.forPlugin(pluginId, deps); } - - /** - * Method to be called during shutdown to destroy all known connections. - */ - async shutdown(deps: { logger: LoggerService }): Promise { - return this.impl.shutdown(deps); - } } diff --git a/packages/backend-defaults/src/entrypoints/database/databaseServiceFactory.ts b/packages/backend-defaults/src/entrypoints/database/databaseServiceFactory.ts index 5e23e6ac4c..525830471b 100644 --- a/packages/backend-defaults/src/entrypoints/database/databaseServiceFactory.ts +++ b/packages/backend-defaults/src/entrypoints/database/databaseServiceFactory.ts @@ -41,21 +41,16 @@ export const databaseServiceFactory = createServiceFactory({ rootLogger: coreServices.rootLogger, }, async createRootContext({ config, rootLifecycle, rootLogger }) { - const databaseManager = config.getOptional('backend.database') - ? DatabaseManager.fromConfig(config) + return config.getOptional('backend.database') + ? DatabaseManager.fromConfig(config, { rootLifecycle, rootLogger }) : DatabaseManager.fromConfig( new ConfigReader({ backend: { database: { client: 'better-sqlite3', connection: ':memory:' }, }, }), + { rootLifecycle, rootLogger }, ); - - rootLifecycle.addShutdownHook(async () => { - await databaseManager.shutdown({ logger: rootLogger }); - }); - - return databaseManager; }, async factory({ pluginMetadata, lifecycle, logger }, databaseManager) { return databaseManager.forPlugin(pluginMetadata.getId(), { From 2bc52ff8bb9ba032971baa0bdbf054b909b30415 Mon Sep 17 00:00:00 2001 From: Eric Peterson Date: Tue, 1 Oct 2024 16:13:22 +0200 Subject: [PATCH 5/9] Update API report for backend defaults Signed-off-by: Eric Peterson --- packages/backend-defaults/report-database.api.md | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/packages/backend-defaults/report-database.api.md b/packages/backend-defaults/report-database.api.md index d9e9e9e287..9c07338264 100644 --- a/packages/backend-defaults/report-database.api.md +++ b/packages/backend-defaults/report-database.api.md @@ -7,6 +7,8 @@ import { DatabaseService } from '@backstage/backend-plugin-api'; import { LifecycleService } from '@backstage/backend-plugin-api'; import { LoggerService } from '@backstage/backend-plugin-api'; import { RootConfigService } from '@backstage/backend-plugin-api'; +import { RootLifecycleService } from '@backstage/backend-plugin-api'; +import { RootLoggerService } from '@backstage/backend-plugin-api'; import { ServiceFactory } from '@backstage/backend-plugin-api'; // @public @@ -27,6 +29,8 @@ export class DatabaseManager { // @public export type DatabaseManagerOptions = { migrations?: DatabaseService['migrations']; + rootLogger?: RootLoggerService; + rootLifecycle?: RootLifecycleService; }; // @public From ffd1f4ab74707b3c76177bae2781dbd2beb29759 Mon Sep 17 00:00:00 2001 From: Eric Peterson Date: Tue, 1 Oct 2024 16:20:18 +0200 Subject: [PATCH 6/9] Add changesets Signed-off-by: Eric Peterson --- .changeset/crash-loop-baby.md | 6 ++++++ .changeset/crash-loop-honey.md | 5 +++++ 2 files changed, 11 insertions(+) create mode 100644 .changeset/crash-loop-baby.md create mode 100644 .changeset/crash-loop-honey.md diff --git a/.changeset/crash-loop-baby.md b/.changeset/crash-loop-baby.md new file mode 100644 index 0000000000..4f629c7f92 --- /dev/null +++ b/.changeset/crash-loop-baby.md @@ -0,0 +1,6 @@ +--- +'@backstage/backend-app-api': patch +'@backstage/backend-defaults': patch +--- + +Plugin lifecycle shutdown hooks are now performed before root lifecycle shutdown hooks. diff --git a/.changeset/crash-loop-honey.md b/.changeset/crash-loop-honey.md new file mode 100644 index 0000000000..60f9f04600 --- /dev/null +++ b/.changeset/crash-loop-honey.md @@ -0,0 +1,5 @@ +--- +'@backstage/backend-defaults': patch +--- + +The database manager now attempts to close any database connections in a root lifecycle shutdown hook. From 8374f0ca17e00d93115b55db0f09947624aac294 Mon Sep 17 00:00:00 2001 From: Eric Peterson Date: Tue, 1 Oct 2024 16:43:03 +0200 Subject: [PATCH 7/9] Test the database manager shutdown hook Signed-off-by: Eric Peterson --- .../database/DatabaseManager.test.ts | 31 +++++++++++++++++++ 1 file changed, 31 insertions(+) diff --git a/packages/backend-defaults/src/entrypoints/database/DatabaseManager.test.ts b/packages/backend-defaults/src/entrypoints/database/DatabaseManager.test.ts index 26c2d54231..c7e739c7c7 100644 --- a/packages/backend-defaults/src/entrypoints/database/DatabaseManager.test.ts +++ b/packages/backend-defaults/src/entrypoints/database/DatabaseManager.test.ts @@ -176,4 +176,35 @@ describe('DatabaseManagerImpl', () => { skip: true, }); }); + + it('registers a shutdown hook if root lifecycle service is provided', async () => { + // Given a database manager that is provided a rootLifecycle service + const rootLifecycle = { addShutdownHook: jest.fn() } as unknown as any; + const destroy = jest.fn(); + const connector1 = { + getClient: jest.fn().mockResolvedValue({ destroy }), + } satisfies Connector; + const impl = new DatabaseManagerImpl( + new ConfigReader({ + client: 'pg', + }), + { + pg: connector1, + }, + { rootLifecycle }, + ); + + // Then a shutdown hook should have been added + expect(rootLifecycle.addShutdownHook).toHaveBeenCalled(); + const shutdownHook = rootLifecycle.addShutdownHook.mock.calls[0][0]; + + // When a database client for a plugin is retrieved + await impl.forPlugin('plugin1', deps).getClient(); + + // And the shutdownhook is called + await shutdownHook(); + + // Then the destroy method should have been called on the resolved client + expect(destroy).toHaveBeenCalled(); + }); }); From 8b49cc2be94ddd05d7fff1aa938a30c322556a58 Mon Sep 17 00:00:00 2001 From: Eric Peterson Date: Wed, 2 Oct 2024 13:43:58 +0200 Subject: [PATCH 8/9] Avoid unnecessary keepalive failed warnings during shutdown Signed-off-by: Eric Peterson --- .../entrypoints/database/DatabaseManager.ts | 48 +++++++++++-------- 1 file changed, 29 insertions(+), 19 deletions(-) diff --git a/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts b/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts index 9441c8bca9..0f6207ef07 100644 --- a/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts +++ b/packages/backend-defaults/src/entrypoints/database/DatabaseManager.ts @@ -57,6 +57,10 @@ export class DatabaseManagerImpl { private readonly connectors: Record, private readonly options?: DatabaseManagerOptions, private readonly databaseCache: Map> = new Map(), + private readonly keepaliveIntervals: Map< + string, + NodeJS.Timeout + > = new Map(), ) { // If a rootLifecycle service was provided, register a shutdown hook to // clean up any database connections. @@ -108,6 +112,9 @@ export class DatabaseManagerImpl { const pluginIds = Array.from(this.databaseCache.keys()); await Promise.allSettled( pluginIds.map(async pluginId => { + // We no longer need to keep connections alive. + clearInterval(this.keepaliveIntervals.get(pluginId)); + const connection = await this.databaseCache.get(pluginId); if (connection) { await connection.destroy().catch((error: unknown) => { @@ -187,25 +194,28 @@ export class DatabaseManagerImpl { ): void { let lastKeepaliveFailed = false; - setInterval(() => { - // During testing it can happen that the environment is torn down and - // this client is `undefined`, but this interval is still run. - client?.raw('select 1').then( - () => { - lastKeepaliveFailed = false; - }, - (error: unknown) => { - if (!lastKeepaliveFailed) { - lastKeepaliveFailed = true; - logger.warn( - `Database keepalive failed for plugin ${pluginId}, ${stringifyError( - error, - )}`, - ); - } - }, - ); - }, 60 * 1000); + this.keepaliveIntervals.set( + pluginId, + setInterval(() => { + // During testing it can happen that the environment is torn down and + // this client is `undefined`, but this interval is still run. + client?.raw('select 1').then( + () => { + lastKeepaliveFailed = false; + }, + (error: unknown) => { + if (!lastKeepaliveFailed) { + lastKeepaliveFailed = true; + logger.warn( + `Database keepalive failed for plugin ${pluginId}, ${stringifyError( + error, + )}`, + ); + } + }, + ); + }, 60 * 1000), + ); } } From e36d12f5badc0ec7277a7009e92dde739e3c8a29 Mon Sep 17 00:00:00 2001 From: Eric Peterson Date: Wed, 2 Oct 2024 14:55:50 +0200 Subject: [PATCH 9/9] Attempt to abort tasks when root lifecycle shutdown hook is invoked Signed-off-by: Eric Peterson --- .changeset/crash-loop-yeah.md | 5 ++ .../backend-defaults/report-database.api.md | 2 + .../backend-defaults/report-scheduler.api.md | 2 + .../scheduler/lib/DefaultSchedulerService.ts | 13 ++++- .../scheduler/lib/LocalTaskWorker.ts | 18 +++--- .../lib/PluginTaskSchedulerImpl.test.ts | 57 +++++++++++++++++++ .../scheduler/lib/PluginTaskSchedulerImpl.ts | 17 +++++- .../entrypoints/scheduler/lib/TaskWorker.ts | 12 ++-- .../scheduler/schedulerServiceFactory.ts | 5 +- 9 files changed, 109 insertions(+), 22 deletions(-) create mode 100644 .changeset/crash-loop-yeah.md diff --git a/.changeset/crash-loop-yeah.md b/.changeset/crash-loop-yeah.md new file mode 100644 index 0000000000..098bc01bae --- /dev/null +++ b/.changeset/crash-loop-yeah.md @@ -0,0 +1,5 @@ +--- +'@backstage/backend-defaults': patch +--- + +The task scheduler now attempts to abort any tasks if it detects that Backstage is being shut down. diff --git a/packages/backend-defaults/report-database.api.md b/packages/backend-defaults/report-database.api.md index 9c07338264..05c5ab251c 100644 --- a/packages/backend-defaults/report-database.api.md +++ b/packages/backend-defaults/report-database.api.md @@ -3,6 +3,8 @@ > Do not edit this file. It is a report generated by [API Extractor](https://api-extractor.com/). ```ts +/// + import { DatabaseService } from '@backstage/backend-plugin-api'; import { LifecycleService } from '@backstage/backend-plugin-api'; import { LoggerService } from '@backstage/backend-plugin-api'; diff --git a/packages/backend-defaults/report-scheduler.api.md b/packages/backend-defaults/report-scheduler.api.md index 351873ad31..7d886e539c 100644 --- a/packages/backend-defaults/report-scheduler.api.md +++ b/packages/backend-defaults/report-scheduler.api.md @@ -5,6 +5,7 @@ ```ts import { DatabaseService } from '@backstage/backend-plugin-api'; import { LoggerService } from '@backstage/backend-plugin-api'; +import { RootLifecycleService } from '@backstage/backend-plugin-api'; import { SchedulerService } from '@backstage/backend-plugin-api'; import { ServiceFactory } from '@backstage/backend-plugin-api'; @@ -14,6 +15,7 @@ export class DefaultSchedulerService { static create(options: { database: DatabaseService; logger: LoggerService; + rootLifecycle?: RootLifecycleService; }): SchedulerService; } diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/DefaultSchedulerService.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/DefaultSchedulerService.ts index 8e2ac5cd07..dcaf673753 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/DefaultSchedulerService.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/DefaultSchedulerService.ts @@ -17,6 +17,7 @@ import { DatabaseService, LoggerService, + RootLifecycleService, SchedulerService, } from '@backstage/backend-plugin-api'; import { once } from 'lodash'; @@ -34,6 +35,7 @@ export class DefaultSchedulerService { static create(options: { database: DatabaseService; logger: LoggerService; + rootLifecycle?: RootLifecycleService; }): SchedulerService { const databaseFactory = once(async () => { const knex = await options.database.getClient(); @@ -43,17 +45,24 @@ export class DefaultSchedulerService { } if (process.env.NODE_ENV !== 'test') { + const abortController = new AbortController(); const janitor = new PluginTaskSchedulerJanitor({ knex, waitBetweenRuns: Duration.fromObject({ minutes: 1 }), logger: options.logger, }); - janitor.start(); + + options.rootLifecycle?.addShutdownHook(() => abortController.abort()); + janitor.start(abortController.signal); } return knex; }); - return new PluginTaskSchedulerImpl(databaseFactory, options.logger); + return new PluginTaskSchedulerImpl( + databaseFactory, + options.logger, + options.rootLifecycle, + ); } } diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/LocalTaskWorker.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/LocalTaskWorker.ts index 3510a855e9..278afae224 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/LocalTaskWorker.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/LocalTaskWorker.ts @@ -36,7 +36,7 @@ export class LocalTaskWorker { private readonly logger: LoggerService, ) {} - start(settings: TaskSettingsV2, options?: { signal?: AbortSignal }) { + start(settings: TaskSettingsV2, options: { signal: AbortSignal }) { this.logger.info( `Task worker starting: ${this.taskId}, ${JSON.stringify(settings)}`, ); @@ -48,18 +48,18 @@ export class LocalTaskWorker { if (settings.initialDelayDuration) { await this.sleep( Duration.fromISO(settings.initialDelayDuration), - options?.signal, + options.signal, ); } - while (!options?.signal?.aborted) { + while (!options.signal.aborted) { const startTime = process.hrtime(); - await this.runOnce(settings, options?.signal); + await this.runOnce(settings, options.signal); const timeTaken = process.hrtime(startTime); await this.waitUntilNext( settings, (timeTaken[0] + timeTaken[1] / 1e9) * 1000, - options?.signal, + options.signal, ); } @@ -89,7 +89,7 @@ export class LocalTaskWorker { */ private async runOnce( settings: TaskSettingsV2, - signal?: AbortSignal, + signal: AbortSignal, ): Promise { // Abort the task execution either if the worker is stopped, or if the // task timeout is hit @@ -115,9 +115,9 @@ export class LocalTaskWorker { private async waitUntilNext( settings: TaskSettingsV2, lastRunMillis: number, - signal?: AbortSignal, + signal: AbortSignal, ) { - if (signal?.aborted) { + if (signal.aborted) { return; } @@ -145,7 +145,7 @@ export class LocalTaskWorker { private async sleep( duration: Duration, - abortSignal?: AbortSignal, + abortSignal: AbortSignal, ): Promise { this.abortWait = delegateAbortController(abortSignal); await sleep(duration, this.abortWait.signal); diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.test.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.test.ts index 3e0782d4a3..3cb41a0b48 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.test.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.test.ts @@ -38,6 +38,7 @@ function defer() { jest.setTimeout(60_000); describe('PluginTaskManagerImpl', () => { + const addShutdownHook = jest.fn(); const databases = TestDatabases.create({ ids: ['POSTGRES_16', 'POSTGRES_12', 'SQLITE_3'], }); @@ -51,12 +52,17 @@ describe('PluginTaskManagerImpl', () => { jest.useFakeTimers(); }, 60_000); + beforeEach(() => { + jest.clearAllMocks(); + }); + async function init(databaseId: TestDatabaseId) { const knex = await databases.init(databaseId); await migrateBackendTasks(knex); const manager = new PluginTaskSchedulerImpl( async () => knex, mockServices.logger.mock(), + { addShutdownHook, addStartupHook: jest.fn() }, ); return { knex, manager }; } @@ -103,6 +109,33 @@ describe('PluginTaskManagerImpl', () => { expect(fn).toHaveBeenCalledWith(expect.any(AbortSignal)); }, ); + + it.each(databases.eachSupportedId())( + 'aborts the task if shutdown hook is invoked, %p', + async databaseId => { + const { manager } = await init(databaseId); + + const fn = jest.fn(); + const promise = new Promise(resolve => + fn.mockImplementation(resolve), + ); + await manager.scheduleTask({ + id: 'task3', + timeout: Duration.fromMillis(5000), + frequency: { cron: '* * * * * *' }, + fn, + scope: 'global', + }); + + const shutdownHook = addShutdownHook.mock.calls[0][0]; + const abortSignal = await promise; + expect(abortSignal.aborted).toBe(false); + + // Should be aborted after the shutdown hook is invoked + await shutdownHook(); + expect(abortSignal.aborted).toBe(true); + }, + ); }); describe('triggerTask with global scope', () => { @@ -212,6 +245,30 @@ describe('PluginTaskManagerImpl', () => { await promise; expect(fn).toHaveBeenCalledWith(expect.any(AbortSignal)); }, 60_000); + + it('aborts the task if shutdown hook is invoked', async () => { + const { manager } = await init('SQLITE_3'); + + const fn = jest.fn(); + const promise = new Promise(resolve => + fn.mockImplementation(resolve), + ); + await manager.scheduleTask({ + id: 'task3', + timeout: Duration.fromMillis(5000), + frequency: { cron: '* * * * * *' }, + fn, + scope: 'local', + }); + + const shutdownHook = addShutdownHook.mock.calls[0][0]; + const abortSignal = await promise; + expect(abortSignal.aborted).toBe(false); + + // Should be aborted after the shutdown hook is invoked + await shutdownHook(); + expect(abortSignal.aborted).toBe(true); + }, 60_000); }); describe('triggerTask with local scope', () => { diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.ts index 5c3af24ab1..5c04521149 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.ts @@ -16,6 +16,7 @@ import { LoggerService, + RootLifecycleService, SchedulerService, SchedulerServiceTaskDescriptor, SchedulerServiceTaskFunction, @@ -29,7 +30,7 @@ import { Duration } from 'luxon'; import { LocalTaskWorker } from './LocalTaskWorker'; import { TaskWorker } from './TaskWorker'; import { TaskSettingsV2 } from './types'; -import { TRACER_ID, validateId } from './util'; +import { delegateAbortController, TRACER_ID, validateId } from './util'; const tracer = trace.getTracer(TRACER_ID); @@ -39,6 +40,7 @@ const tracer = trace.getTracer(TRACER_ID); export class PluginTaskSchedulerImpl implements SchedulerService { private readonly localTasksById = new Map(); private readonly allScheduledTasks: SchedulerServiceTaskDescriptor[] = []; + private readonly shutdownInitiated: Promise; private readonly counter: Counter; private readonly duration: Histogram; @@ -46,6 +48,7 @@ export class PluginTaskSchedulerImpl implements SchedulerService { constructor( private readonly databaseFactory: () => Promise, private readonly logger: LoggerService, + rootLifecycle?: RootLifecycleService, ) { const meter = metrics.getMeter('default'); this.counter = meter.createCounter('backend_tasks.task.runs.count', { @@ -55,6 +58,9 @@ export class PluginTaskSchedulerImpl implements SchedulerService { description: 'Histogram of task run durations', unit: 'seconds', }); + this.shutdownInitiated = new Promise(shutdownInitiated => { + rootLifecycle?.addShutdownHook(() => shutdownInitiated(true)); + }); } async triggerTask(id: string): Promise { @@ -83,6 +89,11 @@ export class PluginTaskSchedulerImpl implements SchedulerService { timeoutAfterDuration: parseDuration(task.timeout), }; + // Delegated abort controller that will abort either when the provided + // controller aborts, or when a root lifecycle shutdown happens + const abortController = delegateAbortController(task.signal); + this.shutdownInitiated.then(() => abortController.abort()); + if (scope === 'global') { const knex = await this.databaseFactory(); const worker = new TaskWorker( @@ -91,14 +102,14 @@ export class PluginTaskSchedulerImpl implements SchedulerService { knex, this.logger.child({ task: task.id }), ); - await worker.start(settings, { signal: task.signal }); + await worker.start(settings, { signal: abortController.signal }); } else { const worker = new LocalTaskWorker( task.id, this.instrumentedFunction(task, scope), this.logger.child({ task: task.id }), ); - worker.start(settings, { signal: task.signal }); + worker.start(settings, { signal: abortController.signal }); this.localTasksById.set(task.id, worker); } diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts index c5570a8e30..48bc81efbf 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts @@ -41,7 +41,7 @@ export class TaskWorker { private readonly workCheckFrequency: Duration = DEFAULT_WORK_CHECK_FREQUENCY, ) {} - async start(settings: TaskSettingsV2, options?: { signal?: AbortSignal }) { + async start(settings: TaskSettingsV2, options: { signal: AbortSignal }) { try { await this.persistTask(settings); } catch (e) { @@ -68,18 +68,18 @@ export class TaskWorker { if (settings.initialDelayDuration) { await sleep( Duration.fromISO(settings.initialDelayDuration), - options?.signal, + options.signal, ); } - while (!options?.signal?.aborted) { - const runResult = await this.runOnce(options?.signal); + while (!options.signal.aborted) { + const runResult = await this.runOnce(options.signal); if (runResult.result === 'abort') { break; } - await sleep(workCheckFrequency, options?.signal); + await sleep(workCheckFrequency, options.signal); } this.logger.info(`Task worker finished: ${this.taskId}`); @@ -122,7 +122,7 @@ export class TaskWorker { * @returns The outcome of the attempt */ private async runOnce( - signal?: AbortSignal, + signal: AbortSignal, ): Promise< | { result: 'not-ready-yet' } | { result: 'abort' } diff --git a/packages/backend-defaults/src/entrypoints/scheduler/schedulerServiceFactory.ts b/packages/backend-defaults/src/entrypoints/scheduler/schedulerServiceFactory.ts index effa5c5925..3668dbd17a 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/schedulerServiceFactory.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/schedulerServiceFactory.ts @@ -34,8 +34,9 @@ export const schedulerServiceFactory = createServiceFactory({ deps: { database: coreServices.database, logger: coreServices.logger, + rootLifecycle: coreServices.rootLifecycle, }, - async factory({ database, logger }) { - return DefaultSchedulerService.create({ database, logger }); + async factory({ database, logger, rootLifecycle }) { + return DefaultSchedulerService.create({ database, logger, rootLifecycle }); }, });