Fix divergent branch issue

Signed-off-by: Damon Kaswell <damon.kaswell1@hp.com>
This commit is contained in:
Damon Kaswell
2022-11-03 08:25:25 -07:00
committed by Fredrik Adelöw
parent 90c1da9277
commit cfe3612e08
2 changed files with 114 additions and 123 deletions
@@ -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<IngestionUpsertIFace>
* @param provider string
* @param update Partial<IngestionUpsertIFace>
*/
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<IngestionRecord>('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<IngestionUpsertIFace> = {
@@ -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<IngestionUpsertIFace>) {
await this.updateIngestionRecordByProvider(provider, update);
}
async insertRecord(options: IngestionRecordInsert) {
await this.insertIngestionRecord(options);
}
}
@@ -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<TCursor, TContext> {
/**
@@ -60,14 +58,20 @@ export interface IncrementalEntityProvider<TCursor, TContext> {
*/
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<TCursor, TContext> {
}
/**
* 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<T> {
/**
* Indicates whether there are any further pages of entities to
* ingest after this one.
*
*/
done: boolean;
@@ -110,7 +111,6 @@ export interface EntityIteratorResult<T> {
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: {