diff --git a/packages/backend-common/api-report.md b/packages/backend-common/api-report.md index 93c25bca24..480b6419fa 100644 --- a/packages/backend-common/api-report.md +++ b/packages/backend-common/api-report.md @@ -375,6 +375,9 @@ export class GitlabUrlReader implements UrlReader { export { isChildPath }; +// @public +export function isDatabaseConflictError(e: unknown): boolean; + // @public export function loadBackendConfig(options: { logger: Logger_2; diff --git a/packages/backend-common/package.json b/packages/backend-common/package.json index a2255926e9..e0d4be1242 100644 --- a/packages/backend-common/package.json +++ b/packages/backend-common/package.json @@ -63,7 +63,6 @@ "minimatch": "^3.0.4", "minimist": "^1.2.5", "morgan": "^1.10.0", - "node-abort-controller": "^3.0.0", "raw-body": "^2.4.1", "selfsigned": "^1.10.7", "stoppable": "^1.1.0", diff --git a/packages/backend-common/src/database/index.ts b/packages/backend-common/src/database/index.ts index 7fcb8bf930..2a820f7367 100644 --- a/packages/backend-common/src/database/index.ts +++ b/packages/backend-common/src/database/index.ts @@ -28,3 +28,4 @@ export { } from './connection'; export type { PluginDatabaseManager } from './types'; +export { isDatabaseConflictError } from './util'; diff --git a/packages/backend-common/src/database/migrateBackendCommon.ts b/packages/backend-common/src/database/migrateBackendCommon.ts index 6283f3a96c..104e757e0b 100644 --- a/packages/backend-common/src/database/migrateBackendCommon.ts +++ b/packages/backend-common/src/database/migrateBackendCommon.ts @@ -16,6 +16,7 @@ import { Knex } from 'knex'; import { resolvePackagePath } from '../paths'; +import { DB_MIGRATIONS_TABLE } from './tables'; const migrationsDir = resolvePackagePath( '@backstage/backend-common', @@ -25,6 +26,6 @@ const migrationsDir = resolvePackagePath( export async function migrateBackendCommon(knex: Knex): Promise { await knex.migrate.latest({ directory: migrationsDir, - tableName: 'knex_migrations_backstage_backend_common', + tableName: DB_MIGRATIONS_TABLE, }); } diff --git a/packages/backend-common/src/database/tables.ts b/packages/backend-common/src/database/tables.ts index e566272ea6..90e8aeca40 100644 --- a/packages/backend-common/src/database/tables.ts +++ b/packages/backend-common/src/database/tables.ts @@ -14,6 +14,7 @@ * limitations under the License. */ +export const DB_MIGRATIONS_TABLE = 'backstage_backend_common__knex_migrations'; export const DB_MUTEXES_TABLE = 'backstage_backend_common__mutexes'; export const DB_TASKS_TABLE = 'backstage_backend_common__tasks'; diff --git a/packages/backend-common/src/database/util.ts b/packages/backend-common/src/database/util.ts new file mode 100644 index 0000000000..607554f294 --- /dev/null +++ b/packages/backend-common/src/database/util.ts @@ -0,0 +1,33 @@ +/* + * Copyright 2021 The Backstage Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +/** + * Tries to deduce whether a thrown error is a database conflict. + * + * @public + * @param e - A thrown error + * @returns True if the error looks like it was a conflict error thrown by a + * known database engine + */ +export function isDatabaseConflictError(e: unknown) { + const message = (e as any)?.message; + + return ( + typeof message === 'string' && + (/SQLITE_CONSTRAINT: UNIQUE/.test(message) || + /unique constraint/.test(message)) + ); +} diff --git a/packages/backend-common/src/tasks/PluginTaskManagerImpl.test.ts b/packages/backend-common/src/tasks/PluginTaskManagerImpl.test.ts index ed5bf0911e..ef5f6db96f 100644 --- a/packages/backend-common/src/tasks/PluginTaskManagerImpl.test.ts +++ b/packages/backend-common/src/tasks/PluginTaskManagerImpl.test.ts @@ -14,8 +14,9 @@ * limitations under the License. */ -import { TestDatabases } from '@backstage/backend-test-utils'; +import { TestDatabaseId, TestDatabases } from '@backstage/backend-test-utils'; import { Duration } from 'luxon'; +import waitForExpect from 'wait-for-expect'; import { migrateBackendCommon } from '../database/migrateBackendCommon'; import { getVoidLogger } from '../logging'; import { PluginTaskManagerImpl } from './PluginTaskManagerImpl'; @@ -25,17 +26,21 @@ describe('PluginTaskManagerImpl', () => { ids: ['POSTGRES_13', 'POSTGRES_9', 'SQLITE_3'], }); - describe('locking', () => { + async function init(databaseId: TestDatabaseId) { + const knex = await databases.init(databaseId); + await migrateBackendCommon(knex); + const manager = new PluginTaskManagerImpl( + async () => knex, + getVoidLogger(), + ); + return { knex, manager }; + } + + describe('acquireLock', () => { it.each(databases.eachSupportedId())( 'can run the happy path, %p', async databaseId => { - const knex = await databases.init(databaseId); - await migrateBackendCommon(knex); - - const manager = new PluginTaskManagerImpl( - async () => knex, - getVoidLogger(), - ); + const { manager } = await init(databaseId); const lock1 = await manager.acquireLock('lock1', { timeout: Duration.fromMillis(5000), @@ -63,5 +68,71 @@ describe('PluginTaskManagerImpl', () => { await (lock1Again as any).release(); }, ); + + it.each(databases.eachSupportedId())( + 'rejects double lock attempts, %p', + async databaseId => { + const { manager } = await init(databaseId); + + const lock1 = await manager.acquireLock('lock1', { + timeout: Duration.fromMillis(5000), + }); + const lock2 = await manager.acquireLock('lock1', { + timeout: Duration.fromMillis(5000), + }); + + expect(lock1.acquired).toBe(true); + expect(lock2.acquired).toBe(false); + + await (lock1 as any).release(); + + const lock1Again = await manager.acquireLock('lock1', { + timeout: Duration.fromMillis(5000), + }); + expect(lock1Again.acquired).toBe(true); + await (lock1Again as any).release(); + }, + ); + + it.each(databases.eachSupportedId())( + 'times out locks, %p', + async databaseId => { + const { manager } = await init(databaseId); + + const lock1 = await manager.acquireLock('lock1', { + timeout: Duration.fromMillis(200), + }); + + expect(lock1.acquired).toBe(true); + + await new Promise(resolve => setTimeout(resolve, 1000)); + + const lock2 = await manager.acquireLock('lock1', { + timeout: Duration.fromMillis(5000), + }); + expect(lock2.acquired).toBe(true); + await (lock2 as any).release(); + }, + ); + }); + + // This is just to test the wrapper code; most of the actual tests are in + // TaskWorker.test.ts + describe('scheduleTask', () => { + it.each(databases.eachSupportedId())( + 'can run the happy path, %p', + async databaseId => { + const { manager } = await init(databaseId); + + const fn = jest.fn(); + const { unschedule } = await manager.scheduleTask('task1', {}, fn); + + await waitForExpect(() => { + expect(fn).toBeCalled(); + }); + + await unschedule(); + }, + ); }); }); diff --git a/packages/backend-common/src/tasks/PluginTaskManagerImpl.ts b/packages/backend-common/src/tasks/PluginTaskManagerImpl.ts index 8b9aad4388..2d1f192fe6 100644 --- a/packages/backend-common/src/tasks/PluginTaskManagerImpl.ts +++ b/packages/backend-common/src/tasks/PluginTaskManagerImpl.ts @@ -18,8 +18,8 @@ import { Knex } from 'knex'; import { Duration } from 'luxon'; import { v4 as uuid } from 'uuid'; import { Logger } from 'winston'; +import { isDatabaseConflictError } from '../database'; import { DbMutexesRow, DB_MUTEXES_TABLE } from '../database/tables'; -import { PluginTaskManagerJanitor } from './PluginTaskManagerJanitor'; import { TaskWorker } from './TaskWorker'; import { PluginTaskManager } from './types'; import { nowPlus, validateId } from './util'; @@ -28,8 +28,6 @@ import { nowPlus, validateId } from './util'; * Implements the actual task management. */ export class PluginTaskManagerImpl implements PluginTaskManager { - private janitor: PluginTaskManagerJanitor | undefined; - constructor( private readonly databaseFactory: () => Promise, private readonly logger: Logger, @@ -47,8 +45,6 @@ export class PluginTaskManagerImpl implements PluginTaskManager { validateId(id); const knex = await this.databaseFactory(); - await this.ensureJanitor(knex); - const ticket = uuid(); async function release() { @@ -62,30 +58,32 @@ export class PluginTaskManagerImpl implements PluginTaskManager { } } + const record: Knex.DbRecord = { + current_lock_ticket: ticket, + current_lock_acquired_at: knex.fn.now(), + current_lock_expires_at: options.timeout + ? nowPlus(options.timeout, knex) + : knex.raw('null'), + }; + // First try to overwrite an existing lock, that has timed out const stolen = await knex(DB_MUTEXES_TABLE) .where('id', '=', id) .whereNotNull('current_lock_ticket') .where('current_lock_expires_at', '<', knex.fn.now()) - .update({ - current_lock_ticket: ticket, - current_lock_acquired_at: knex.fn.now(), - current_lock_expires_at: nowPlus(options.timeout, knex), - }); + .update(record); if (stolen) { return { acquired: true, release }; } try { - await knex(DB_MUTEXES_TABLE).insert({ - id, - current_lock_ticket: ticket, - current_lock_acquired_at: knex.fn.now(), - current_lock_expires_at: nowPlus(options.timeout, knex), - }); + await knex(DB_MUTEXES_TABLE).insert({ id, ...record }); return { acquired: true, release }; - } catch { + } catch (e) { + if (!isDatabaseConflictError(e)) { + this.logger.warn(`Failed to acquire lock, ${e}`); + } return { acquired: false }; } } @@ -102,7 +100,6 @@ export class PluginTaskManagerImpl implements PluginTaskManager { validateId(id); const knex = await this.databaseFactory(); - await this.ensureJanitor(knex); const task = new TaskWorker(id, fn, knex, this.logger); await task.start({ @@ -118,15 +115,4 @@ export class PluginTaskManagerImpl implements PluginTaskManager { }, }; } - - private async ensureJanitor(knex: Knex) { - if (!this.janitor) { - this.janitor = new PluginTaskManagerJanitor({ - knex, - waitBetweenRuns: Duration.fromObject({ minutes: 1 }), - logger: this.logger, - }); - this.janitor.start(); - } - } } diff --git a/packages/backend-common/src/tasks/PluginTaskManagerJanitor.ts b/packages/backend-common/src/tasks/PluginTaskManagerJanitor.ts new file mode 100644 index 0000000000..26af986025 --- /dev/null +++ b/packages/backend-common/src/tasks/PluginTaskManagerJanitor.ts @@ -0,0 +1,114 @@ +/* + * Copyright 2021 The Backstage Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import { Knex } from 'knex'; +import { Duration } from 'luxon'; +import { Logger } from 'winston'; +import { + DbMutexesRow, + DbTasksRow, + DB_MUTEXES_TABLE, + DB_TASKS_TABLE, +} from '../database/tables'; +import { CancelToken } from './CancelToken'; + +/** + * Makes sure to auto-expire and clean up things that time out or for other + * reasons should not be left lingering. + */ +export class PluginTaskManagerJanitor { + private readonly knex: Knex; + private readonly waitBetweenRuns: Duration; + private readonly logger: Logger; + private readonly cancelToken: CancelToken; + + constructor(options: { + knex: Knex; + waitBetweenRuns: Duration; + logger: Logger; + }) { + this.knex = options.knex; + this.waitBetweenRuns = options.waitBetweenRuns; + this.logger = options.logger; + this.cancelToken = CancelToken.create(); + } + + async start() { + while (!this.cancelToken.isCancelled) { + try { + await this.runOnce(); + } catch (e) { + this.logger.warn(`Error while performing janitorial tasks, ${e}`); + } + + await this.sleep(this.waitBetweenRuns); + } + } + + async stop() { + this.cancelToken.cancel(); + } + + private async runOnce() { + // SQLite currently (Oct 1 2021) returns a number for returning() + // statements, effectively ignoring them and instead returning the outcome + // of the delete() - and knex also emits a warning about that fact, which + // is why we avoid that entirely for the sqlite3 driver. + // https://github.com/knex/knex/issues/4370 + // https://github.com/mapbox/node-sqlite3/issues/1453 + + const mutexesQuery = this.knex(DB_MUTEXES_TABLE) + .where('current_lock_expires_at', '<', this.knex.fn.now()) + .delete(); + + if (this.knex.client.config.client === 'sqlite3') { + const mutexes = await mutexesQuery; + if (mutexes > 0) { + this.logger.warn(`${mutexes} mutex locks timed out and were lost`); + } + } else { + const mutexes = await mutexesQuery.returning(['id']); + for (const { id } of mutexes) { + this.logger.warn(`Mutex lock timed out and was lost: ${id}`); + } + } + + const tasksQuery = this.knex(DB_TASKS_TABLE) + .where('current_run_expires_at', '<', this.knex.fn.now()) + .delete(); + + if (this.knex.client.config.client === 'sqlite3') { + const tasks = await tasksQuery; + this.logger.warn(`${tasks} tasks timed out and were lost`); + } else { + const tasks = await tasksQuery.returning(['id']); + for (const { id } of tasks) { + this.logger.warn(`Task timed out and was lost: ${id}`); + } + } + } + + /** + * Sleeps for the given duration, but aborts sooner if the cancel token + * triggers. + */ + private async sleep(duration: Duration) { + await Promise.race([ + new Promise(resolve => setTimeout(resolve, duration.as('milliseconds'))), + this.cancelToken.promise, + ]); + } +} diff --git a/packages/backend-common/src/tasks/TaskManager.ts b/packages/backend-common/src/tasks/TaskManager.ts index edb8e46609..ea49c2383e 100644 --- a/packages/backend-common/src/tasks/TaskManager.ts +++ b/packages/backend-common/src/tasks/TaskManager.ts @@ -16,11 +16,13 @@ import { Config } from '@backstage/config'; import { memoize } from 'lodash'; +import { Duration } from 'luxon'; import { Logger } from 'winston'; import { DatabaseManager } from '../database'; import { migrateBackendCommon } from '../database/migrateBackendCommon'; import { getRootLogger } from '../logging'; import { PluginTaskManagerImpl } from './PluginTaskManagerImpl'; +import { PluginTaskManagerJanitor } from './PluginTaskManagerJanitor'; import { PluginTaskManager } from './types'; /** @@ -52,7 +54,16 @@ export class TaskManager { forPlugin(pluginId: string): PluginTaskManager { const databaseFactory = memoize(async () => { const knex = await this.databaseManager.forPlugin(pluginId).getClient(); + await migrateBackendCommon(knex); + + const janitor = new PluginTaskManagerJanitor({ + knex, + waitBetweenRuns: Duration.fromObject({ minutes: 1 }), + logger: this.logger, + }); + janitor.start(); + return knex; });