From cfe3612e085423ab16473fd57c67378722147765 Mon Sep 17 00:00:00 2001 From: Damon Kaswell Date: Thu, 3 Nov 2022 08:25:25 -0700 Subject: [PATCH] Fix divergent branch issue Signed-off-by: Damon Kaswell --- .../IncrementalIngestionDatabaseManager.ts | 131 +++++++++--------- .../src/types.ts | 106 +++++++------- 2 files changed, 114 insertions(+), 123 deletions(-) diff --git a/plugins/incremental-ingestion-backend/src/database/IncrementalIngestionDatabaseManager.ts b/plugins/incremental-ingestion-backend/src/database/IncrementalIngestionDatabaseManager.ts index d7dc7dc1de..6a59757ca1 100644 --- a/plugins/incremental-ingestion-backend/src/database/IncrementalIngestionDatabaseManager.ts +++ b/plugins/incremental-ingestion-backend/src/database/IncrementalIngestionDatabaseManager.ts @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + import { Knex } from 'knex'; import type { DeferredEntity } from '@backstage/plugin-catalog-backend'; import { stringifyEntityRef } from '@backstage/catalog-model'; @@ -21,14 +22,12 @@ import { v4 } from 'uuid'; import { INCREMENTAL_ENTITY_PROVIDER_ANNOTATION, IngestionRecord, - IngestionRecordInsert, IngestionRecordUpdate, IngestionUpsertIFace, MarkRecord, MarkRecordInsert, } from '../types'; -/** @public */ export class IncrementalIngestionDatabaseManager { private client: Knex; @@ -38,7 +37,7 @@ export class IncrementalIngestionDatabaseManager { /** * Performs an update to the ingestion record with matching `id`. - * @param options - IngestionRecordUpdate + * @param options IngestionRecordUpdate */ async updateIngestionRecordById(options: IngestionRecordUpdate) { await this.client.transaction(async tx => { @@ -49,8 +48,8 @@ export class IncrementalIngestionDatabaseManager { /** * Performs an update to the ingestion record with matching provider name. Will only update active records. - * @param provider - string - * @param update - Partial + * @param provider string + * @param update Partial */ async updateIngestionRecordByProvider( provider: string, @@ -59,18 +58,17 @@ export class IncrementalIngestionDatabaseManager { await this.client.transaction(async tx => { await tx('ingestion.ingestions') .where('provider_name', provider) - .andWhere('rest_completed_at', null) + .andWhere('completion_ticket', 'open') .update(update); }); } /** * Performs an insert into the `ingestion.ingestions` table with the supplied values. - * @param options - IngestionRecordInsert + * @param record IngestionUpsertIFace */ - async insertIngestionRecord(options: IngestionRecordInsert) { + async insertIngestionRecord(record: IngestionUpsertIFace) { await this.client.transaction(async tx => { - const { record } = options; await tx('ingestion.ingestions').insert(record); }); } @@ -102,14 +100,14 @@ export class IncrementalIngestionDatabaseManager { /** * Finds the current ingestion record for the named provider. - * @param provider - string + * @param provider string * @returns IngestionRecord | undefined */ async getCurrentIngestionRecord(provider: string) { return await this.client.transaction(async tx => { const record = await tx('ingestion.ingestions') .where('provider_name', provider) - .andWhere('rest_completed_at', null) + .andWhere('completion_ticket', 'open') .first(); return record; }); @@ -117,8 +115,8 @@ export class IncrementalIngestionDatabaseManager { /** * Removes all entries from `ingestion_marks_entities`, `ingestion_marks`, and `ingestions` - * for prior ingestions that completed (i.e., have a value for `rest_completed_at`). - * @param provider - string + * for prior ingestions that completed (i.e., have a `completion_ticket` value other than 'open'). + * @param provider string * @returns A count of deletions for each record type. */ async clearFinishedIngestions(provider: string) { @@ -134,7 +132,7 @@ export class IncrementalIngestionDatabaseManager { tx('ingestion.ingestions') .select('id') .where('provider_name', provider) - .andWhereNot('rest_completed_at', null), + .andWhereNot('completion_ticket', 'open'), ), ); @@ -145,13 +143,13 @@ export class IncrementalIngestionDatabaseManager { tx('ingestion.ingestions') .select('id') .where('provider_name', provider) - .andWhereNot('rest_completed_at', null), + .andWhereNot('completion_ticket', 'open'), ); const ingestionsDeleted = await tx('ingestion.ingestions') .delete() .where('provider_name', provider) - .andWhereNot('rest_completed_at', null); + .andWhereNot('completion_ticket', 'open'); return { deletions: { @@ -167,8 +165,8 @@ export class IncrementalIngestionDatabaseManager { * Automatically cleans up duplicate ingestion records if they were accidentally created. * Any ingestion record where the `rest_completed_at` is null (meaning it is active) AND * the ingestionId is incorrect is a duplicate ingestion record. - * @param ingestionId - string - * @param provider - string + * @param ingestionId string + * @param provider string */ async clearDuplicateIngestions(ingestionId: string, provider: string) { await this.client.transaction(async tx => { @@ -197,7 +195,7 @@ export class IncrementalIngestionDatabaseManager { /** * This method fully purges and resets all ingestion records for the named provider, and * leaves it in a paused state. - * @param provider - string + * @param provider string * @returns Counts of all deleted ingestion records */ async purgeAndResetProvider(provider: string) { @@ -248,15 +246,14 @@ export class IncrementalIngestionDatabaseManager { const next_action_at = new Date(); next_action_at.setTime(next_action_at.getTime() + 24 * 60 * 60 * 1000); - await this.insertRecord({ - record: { - id: v4(), - next_action: 'rest', - provider_name: provider, - next_action_at, - ingestion_completed_at: new Date(), - status: 'resting', - }, + await this.insertIngestionRecord({ + id: v4(), + next_action: 'rest', + provider_name: provider, + next_action_at, + ingestion_completed_at: new Date(), + status: 'resting', + completion_ticket: v4(), }); return { provider, ingestionsDeleted, marksDeleted, markEntitiesDeleted }; @@ -265,27 +262,31 @@ export class IncrementalIngestionDatabaseManager { /** * Creates a new ingestion record. - * @param provider - string + * @param provider string * @returns A new ingestion record */ async createProviderIngestionRecord(provider: string) { const ingestionId = v4(); const nextAction = 'ingest'; - await this.insertIngestionRecord({ - record: { + try { + await this.insertIngestionRecord({ id: ingestionId, next_action: nextAction, provider_name: provider, status: 'bursting', - }, - }); - return { ingestionId, nextAction, attempts: 0, nextActionAt: Date.now() }; + completion_ticket: 'open', + }); + return { ingestionId, nextAction, attempts: 0, nextActionAt: Date.now() }; + } catch (_e) { + // Creating the ingestion record failed. Return undefined. + return undefined; + } } /** * Computes which entities to remove, if any, at the end of a burst. - * @param provider - string - * @param ingestionId - string + * @param provider string + * @param ingestionId string * @returns All entities to remove for this burst. */ async computeRemoved(provider: string, ingestionId: string) { @@ -340,7 +341,7 @@ export class IncrementalIngestionDatabaseManager { /** * Skips any wait time for the next action to run. - * @param provider - string + * @param provider string */ async triggerNextProviderAction(provider: string) { await this.updateIngestionRecordByProvider(provider, { @@ -367,14 +368,13 @@ export class IncrementalIngestionDatabaseManager { for (const provider of providers) { await this.insertIngestionRecord({ - record: { - id: v4(), - next_action: 'rest', - provider_name: provider, - next_action_at, - ingestion_completed_at: new Date(), - status: 'resting', - }, + id: v4(), + next_action: 'rest', + provider_name: provider, + next_action_at, + ingestion_completed_at: new Date(), + status: 'resting', + completion_ticket: v4(), }); } @@ -390,7 +390,7 @@ export class IncrementalIngestionDatabaseManager { /** * Configures the current ingestion record to ingest a burst. - * @param ingestionId - string + * @param ingestionId string */ async setProviderIngesting(ingestionId: string) { await this.updateIngestionRecordById({ @@ -401,7 +401,7 @@ export class IncrementalIngestionDatabaseManager { /** * Indicates the provider is currently ingesting a burst. - * @param ingestionId - string + * @param ingestionId string */ async setProviderBursting(ingestionId: string) { await this.updateIngestionRecordById({ @@ -412,7 +412,7 @@ export class IncrementalIngestionDatabaseManager { /** * Finalizes the current ingestion record to indicate that the post-ingestion rest period is complete. - * @param ingestionId - string + * @param ingestionId string */ async setProviderComplete(ingestionId: string) { await this.updateIngestionRecordById({ @@ -421,14 +421,15 @@ export class IncrementalIngestionDatabaseManager { next_action: 'nothing (done)', rest_completed_at: new Date(), status: 'complete', + completion_ticket: v4(), }, }); } /** * Marks ingestion as complete and starts the post-ingestion rest cycle. - * @param ingestionId - string - * @param restLength - Duration + * @param ingestionId string + * @param restLength Duration */ async setProviderResting(ingestionId: string, restLength: Duration) { await this.updateIngestionRecordById({ @@ -444,7 +445,7 @@ export class IncrementalIngestionDatabaseManager { /** * Marks ingestion as paused after a burst completes. - * @param ingestionId - string + * @param ingestionId string */ async setProviderInterstitial(ingestionId: string) { await this.updateIngestionRecordById({ @@ -455,8 +456,8 @@ export class IncrementalIngestionDatabaseManager { /** * Starts the cancel process for the current ingestion. - * @param ingestionId - string - * @param message - string (optional) + * @param ingestionId string + * @param message string (optional) */ async setProviderCanceling(ingestionId: string, message?: string) { const update: Partial = { @@ -470,7 +471,7 @@ export class IncrementalIngestionDatabaseManager { /** * Completes the cancel process and triggers a new ingestion. - * @param ingestionId - string + * @param ingestionId string */ async setProviderCanceled(ingestionId: string) { await this.updateIngestionRecordById({ @@ -479,16 +480,17 @@ export class IncrementalIngestionDatabaseManager { next_action: 'nothing (canceled)', rest_completed_at: new Date(), status: 'complete', + completion_ticket: v4(), }, }); } /** * Configures the current ingestion to wait and retry, due to a data source error. - * @param ingestionId - string - * @param attempts - number - * @param error - Error - * @param backoffLength - number + * @param ingestionId string + * @param attempts number + * @param error Error + * @param backoffLength number */ async setProviderBackoff( ingestionId: string, @@ -510,7 +512,7 @@ export class IncrementalIngestionDatabaseManager { /** * Returns the last record from `ingestion.ingestion_marks` for the supplied ingestionId. - * @param ingestionId - string + * @param ingestionId string * @returns MarkRecord | undefined */ async getLastMark(ingestionId: string) { @@ -534,7 +536,7 @@ export class IncrementalIngestionDatabaseManager { /** * Performs an insert into the `ingestion.ingestion_marks` table with the supplied values. - * @param options - MarkRecordInsert + * @param options MarkRecordInsert */ async createMark(options: MarkRecordInsert) { const { record } = options; @@ -544,8 +546,8 @@ export class IncrementalIngestionDatabaseManager { } /** * Performs an upsert to the `ingestion.ingestion_mark_entities` table for all deferred entities. - * @param markId - string - * @param entities - DeferredEntity[] + * @param markId string + * @param entities DeferredEntity[] */ async createMarkEntities(markId: string, entities: DeferredEntity[]) { await this.client.transaction(async tx => { @@ -561,7 +563,7 @@ export class IncrementalIngestionDatabaseManager { /** * Deletes the entire content of a table, and returns the number of records deleted. - * @param table - string + * @param table string * @returns number */ async purgeTable(table: string) { @@ -586,7 +588,4 @@ export class IncrementalIngestionDatabaseManager { async updateByName(provider: string, update: Partial) { await this.updateIngestionRecordByProvider(provider, update); } - async insertRecord(options: IngestionRecordInsert) { - await this.insertIngestionRecord(options); - } } diff --git a/plugins/incremental-ingestion-backend/src/types.ts b/plugins/incremental-ingestion-backend/src/types.ts index c9f3f50e1d..5cc0a5072d 100644 --- a/plugins/incremental-ingestion-backend/src/types.ts +++ b/plugins/incremental-ingestion-backend/src/types.ts @@ -30,7 +30,6 @@ import type { PermissionAuthorizer } from '@backstage/plugin-permission-common'; import type { DurationObjectUnits } from 'luxon'; import type { Logger } from 'winston'; import { IncrementalIngestionDatabaseManager } from './database/IncrementalIngestionDatabaseManager'; - /** * Entity annotation containing the incremental entity provider. * @@ -51,7 +50,6 @@ export const INCREMENTAL_ENTITY_PROVIDER_ANNOTATION = * batches of entities in sequence so that you never need to have more * than a few hundred in memory at a time. * - * @public */ export interface IncrementalEntityProvider { /** @@ -60,14 +58,20 @@ export interface IncrementalEntityProvider { */ getProviderName(): string; + /** + * Return a function to register the end of the metric that was + * initialized before the fisrt burst call. + */ + prometheusRegister?: { markIngestionCompleted: (result: string) => void }; + /** * Return a single page of entities from a specific point in the * ingestion. * * @param context - anything needed in order to fetch a single page. * @param cursor - a uniqiue value identifying the page to ingest. - * @returns the entities to be ingested, as well as the cursor of - * the the next page after this one. + * @returns The entities to be ingested, as well as the cursor of + * the next page after this one. */ next( context: TContext, @@ -85,16 +89,13 @@ export interface IncrementalEntityProvider { } /** - * Value returned by an {@link IncrementalEntityProvider} to provide a + * Value returned by an @{link IncrementalEntityProvider} to provide a * single page of entities to ingest. - * - * @public */ export interface EntityIteratorResult { /** * Indicates whether there are any further pages of entities to * ingest after this one. - * */ done: boolean; @@ -110,7 +111,6 @@ export interface EntityIteratorResult { entities: DeferredEntity[]; } -/** @public */ export interface IncrementalEntityProviderOptions { /** * Entities are ingested in bursts. This interval determines how @@ -134,16 +134,11 @@ export interface IncrementalEntityProviderOptions { /** * In the event of an error during an ingestion burst, the backoff * determines how soon it will be retried. E.g. - * `[{ minutes: 1}, { minutes: 5}, {minutes: 30 }, { hours: 3 }]` + * [{ minutes: 1}, { minutes: 5}, {minutes: 30 }, { hours: 3 }] */ backoff?: DurationObjectUnits[]; } -/** - * Provides the environment used by the incremental entity provider, - * - * @public - */ export type PluginEnvironment = { logger: Logger; database: PluginDatabaseManager; @@ -153,20 +148,10 @@ export type PluginEnvironment = { permissions: PermissionAuthorizer; }; -/** - * The core ingestion engine implements this interface - * - * @public - */ export interface IterationEngine { taskFn: TaskFunction; } -/** - * Options passed to the core ingestion engine during initialization. - * - * @public - */ export interface IterationEngineOptions { logger: Logger; connection: EntityProviderConnection; @@ -179,10 +164,15 @@ export interface IterationEngineOptions { /** * The shape of data inserted into or updated in the `ingestion.ingestions` table. - * - * @public */ export interface IngestionUpsertIFace { + /** + * The ingestion record id. + */ + id?: string; + /** + * The next action the incremental entity provider will take. + */ next_action: | 'rest' | 'ingest' @@ -190,6 +180,9 @@ export interface IngestionUpsertIFace { | 'cancel' | 'nothing (done)' | 'nothing (canceled)'; + /** + * Current status of the incremental entity provider. + */ status: | 'complete' | 'bursting' @@ -197,29 +190,38 @@ export interface IngestionUpsertIFace { | 'canceling' | 'interstitial' | 'backing off'; + /** + * The name of the incremental entity provider being updated. + */ provider_name: string; + /** + * Date/time stamp for when the next action will trigger. + */ next_action_at?: Date; - last_error?: string; + /** + * A record of the last error generated by the incremental entity provider. + */ + last_error?: string | null; + /** + * The number of attempts the provider has attempted during the current cycle. + */ attempts?: number; - ingestion_completed_at?: Date; - rest_completed_at?: Date; -} - -/** - * This interface supplies all potential values that can be inserted into the `ingestion.ingestions` table. - * - * @public - */ -export interface IngestionRecordInsert { - record: IngestionUpsertIFace & { - id: string; - }; + /** + * Date/time stamp for the completion of ingestion. + */ + ingestion_completed_at?: Date | string | null; + /** + * Date/time stamp for the end of the rest cycle before the next ingestion. + */ + rest_completed_at?: Date | string | null; + /** + * A record of the finalized status of the ingestion record. Values are either 'open' or a uuid. + */ + completion_ticket: string; } /** * This interface is for updating an existing ingestion record. - * - * @public */ export interface IngestionRecordUpdate { ingestionId: string; @@ -228,8 +230,6 @@ export interface IngestionRecordUpdate { /** * The expected response from the `ingestion.ingestion_marks` table. - * - * @public */ export interface MarkRecord { id: string; @@ -241,26 +241,18 @@ export interface MarkRecord { /** * The expected response from the `ingestion.ingestions` table. - * - * @public */ -export interface IngestionRecord { +export interface IngestionRecord extends IngestionUpsertIFace { id: string; - provider_name: string; - status: string; - next_action: string; next_action_at: Date; - last_error: string | null; - attempts: number; + /** + * The date/time the ingestion record was created. + */ created_at: string; - rest_completed_at: string | null; - ingestion_completed_at: string | null; } /** * This interface supplies all the values for adding an ingestion mark. - * - * @public */ export interface MarkRecordInsert { record: {