Merge pull request #17534 from mikebryant/m/add-observability-catalog-processors

feat(catalog-backend): Add observability for catalog processing
This commit is contained in:
Fredrik Adelöw
2023-07-28 20:58:32 +02:00
committed by GitHub
5 changed files with 494 additions and 297 deletions
+5
View File
@@ -0,0 +1,5 @@
---
'@backstage/plugin-catalog-backend': minor
---
Added OpenTelemetry spans for catalog processing
@@ -23,7 +23,7 @@ import { assertError, serializeError, stringifyError } from '@backstage/errors';
import { Hash } from 'crypto';
import stableStringify from 'fast-json-stable-stringify';
import { Logger } from 'winston';
import { metrics } from '@opentelemetry/api';
import { metrics, trace } from '@opentelemetry/api';
import { ProcessingDatabase, RefreshStateItem } from '../database/types';
import { createCounterMetric, createSummaryMetric } from '../util/metrics';
import {
@@ -35,9 +35,16 @@ import { Stitcher } from '../stitching/Stitcher';
import { startTaskPipeline } from './TaskPipeline';
import { PluginTaskScheduler } from '@backstage/backend-tasks';
import { Config } from '@backstage/config';
import {
addEntityAttributes,
TRACER_ID,
withActiveSpan,
} from '../util/opentelemetry';
const CACHE_TTL = 5;
const tracer = trace.getTracer(TRACER_ID);
export type ProgressTracker = ReturnType<typeof progressTracker>;
export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine {
@@ -131,177 +138,181 @@ export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine {
}
},
processTask: async item => {
const track = this.tracker.processStart(item, this.logger);
await withActiveSpan(tracer, 'ProcessingRun', async span => {
const track = this.tracker.processStart(item, this.logger);
addEntityAttributes(span, item.unprocessedEntity);
try {
const {
id,
state,
unprocessedEntity,
entityRef,
locationKey,
resultHash: previousResultHash,
} = item;
const result = await this.orchestrator.process({
entity: unprocessedEntity,
state,
});
try {
const {
id,
state,
unprocessedEntity,
entityRef,
locationKey,
resultHash: previousResultHash,
} = item;
const result = await this.orchestrator.process({
entity: unprocessedEntity,
state,
});
track.markProcessorsCompleted(result);
track.markProcessorsCompleted(result);
if (result.ok) {
const { ttl: _, ...stateWithoutTtl } = state ?? {};
if (
stableStringify(stateWithoutTtl) !== stableStringify(result.state)
) {
if (result.ok) {
const { ttl: _, ...stateWithoutTtl } = state ?? {};
if (
stableStringify(stateWithoutTtl) !==
stableStringify(result.state)
) {
await this.processingDatabase.transaction(async tx => {
await this.processingDatabase.updateEntityCache(tx, {
id,
state: {
ttl: CACHE_TTL,
...result.state,
},
});
});
}
} else {
const maybeTtl = state?.ttl;
const ttl = Number.isInteger(maybeTtl) ? (maybeTtl as number) : 0;
await this.processingDatabase.transaction(async tx => {
await this.processingDatabase.updateEntityCache(tx, {
id,
state: {
ttl: CACHE_TTL,
...result.state,
},
state: ttl > 0 ? { ...state, ttl: ttl - 1 } : {},
});
});
}
} else {
const maybeTtl = state?.ttl;
const ttl = Number.isInteger(maybeTtl) ? (maybeTtl as number) : 0;
await this.processingDatabase.transaction(async tx => {
await this.processingDatabase.updateEntityCache(tx, {
id,
state: ttl > 0 ? { ...state, ttl: ttl - 1 } : {},
const location =
unprocessedEntity?.metadata?.annotations?.[ANNOTATION_LOCATION];
for (const error of result.errors) {
this.logger.warn(error.message, {
entity: entityRef,
location,
});
});
}
}
const errorsString = JSON.stringify(
result.errors.map(e => serializeError(e)),
);
const location =
unprocessedEntity?.metadata?.annotations?.[ANNOTATION_LOCATION];
for (const error of result.errors) {
this.logger.warn(error.message, {
entity: entityRef,
location,
});
}
const errorsString = JSON.stringify(
result.errors.map(e => serializeError(e)),
);
let hashBuilder = this.createHash().update(errorsString);
let hashBuilder = this.createHash().update(errorsString);
if (result.ok) {
const { entityRefs: parents } =
await this.processingDatabase.transaction(tx =>
this.processingDatabase.listParents(tx, {
entityRef,
}),
);
hashBuilder = hashBuilder
.update(stableStringify({ ...result.completedEntity }))
.update(stableStringify([...result.deferredEntities]))
.update(stableStringify([...result.relations]))
.update(stableStringify([...result.refreshKeys]))
.update(stableStringify([...parents]));
}
const resultHash = hashBuilder.digest('hex');
if (resultHash === previousResultHash) {
// If nothing changed in our produced outputs, we cannot have any
// significant effect on our surroundings; therefore, we just abort
// without any updates / stitching.
track.markSuccessfulWithNoChanges();
return;
}
// If the result was marked as not OK, it signals that some part of the
// processing pipeline threw an exception. This can happen both as part of
// non-catastrophic things such as due to validation errors, as well as if
// something fatal happens inside the processing for other reasons. In any
// case, this means we can't trust that anything in the output is okay. So
// just store the errors and trigger a stich so that they become visible to
// the outside.
if (!result.ok) {
// notify the error listener if the entity can not be processed.
Promise.resolve(undefined)
.then(() =>
this.onProcessingError?.({
unprocessedEntity,
errors: result.errors,
}),
)
.catch(error => {
this.logger.debug(
`Processing error listener threw an exception, ${stringifyError(
error,
)}`,
if (result.ok) {
const { entityRefs: parents } =
await this.processingDatabase.transaction(tx =>
this.processingDatabase.listParents(tx, {
entityRef,
}),
);
});
hashBuilder = hashBuilder
.update(stableStringify({ ...result.completedEntity }))
.update(stableStringify([...result.deferredEntities]))
.update(stableStringify([...result.relations]))
.update(stableStringify([...result.refreshKeys]))
.update(stableStringify([...parents]));
}
const resultHash = hashBuilder.digest('hex');
if (resultHash === previousResultHash) {
// If nothing changed in our produced outputs, we cannot have any
// significant effect on our surroundings; therefore, we just abort
// without any updates / stitching.
track.markSuccessfulWithNoChanges();
return;
}
// If the result was marked as not OK, it signals that some part of the
// processing pipeline threw an exception. This can happen both as part of
// non-catastrophic things such as due to validation errors, as well as if
// something fatal happens inside the processing for other reasons. In any
// case, this means we can't trust that anything in the output is okay. So
// just store the errors and trigger a stich so that they become visible to
// the outside.
if (!result.ok) {
// notify the error listener if the entity can not be processed.
Promise.resolve(undefined)
.then(() =>
this.onProcessingError?.({
unprocessedEntity,
errors: result.errors,
}),
)
.catch(error => {
this.logger.debug(
`Processing error listener threw an exception, ${stringifyError(
error,
)}`,
);
});
await this.processingDatabase.transaction(async tx => {
await this.processingDatabase.updateProcessedEntityErrors(tx, {
id,
errors: errorsString,
resultHash,
});
});
await this.stitcher.stitch(
new Set([stringifyEntityRef(unprocessedEntity)]),
);
track.markSuccessfulWithErrors();
return;
}
result.completedEntity.metadata.uid = id;
let oldRelationSources: Map<string, string>;
await this.processingDatabase.transaction(async tx => {
await this.processingDatabase.updateProcessedEntityErrors(tx, {
id,
errors: errorsString,
resultHash,
});
const { previous } =
await this.processingDatabase.updateProcessedEntity(tx, {
id,
processedEntity: result.completedEntity,
resultHash,
errors: errorsString,
relations: result.relations,
deferredEntities: result.deferredEntities,
locationKey,
refreshKeys: result.refreshKeys,
});
oldRelationSources = new Map(
previous.relations.map(r => [
`${r.source_entity_ref}:${r.type}`,
r.source_entity_ref,
]),
);
});
await this.stitcher.stitch(
new Set([stringifyEntityRef(unprocessedEntity)]),
const newRelationSources = new Map<string, string>(
result.relations.map(relation => {
const sourceEntityRef = stringifyEntityRef(relation.source);
return [`${sourceEntityRef}:${relation.type}`, sourceEntityRef];
}),
);
track.markSuccessfulWithErrors();
return;
const setOfThingsToStitch = new Set<string>([
stringifyEntityRef(result.completedEntity),
]);
newRelationSources.forEach((sourceEntityRef, uniqueKey) => {
if (!oldRelationSources.has(uniqueKey)) {
setOfThingsToStitch.add(sourceEntityRef);
}
});
oldRelationSources!.forEach((sourceEntityRef, uniqueKey) => {
if (!newRelationSources.has(uniqueKey)) {
setOfThingsToStitch.add(sourceEntityRef);
}
});
await this.stitcher.stitch(setOfThingsToStitch);
track.markSuccessfulWithChanges(setOfThingsToStitch.size);
} catch (error) {
assertError(error);
track.markFailed(error);
}
result.completedEntity.metadata.uid = id;
let oldRelationSources: Map<string, string>;
await this.processingDatabase.transaction(async tx => {
const { previous } =
await this.processingDatabase.updateProcessedEntity(tx, {
id,
processedEntity: result.completedEntity,
resultHash,
errors: errorsString,
relations: result.relations,
deferredEntities: result.deferredEntities,
locationKey,
refreshKeys: result.refreshKeys,
});
oldRelationSources = new Map(
previous.relations.map(r => [
`${r.source_entity_ref}:${r.type}`,
r.source_entity_ref,
]),
);
});
const newRelationSources = new Map<string, string>(
result.relations.map(relation => {
const sourceEntityRef = stringifyEntityRef(relation.source);
return [`${sourceEntityRef}:${relation.type}`, sourceEntityRef];
}),
);
const setOfThingsToStitch = new Set<string>([
stringifyEntityRef(result.completedEntity),
]);
newRelationSources.forEach((sourceEntityRef, uniqueKey) => {
if (!oldRelationSources.has(uniqueKey)) {
setOfThingsToStitch.add(sourceEntityRef);
}
});
oldRelationSources!.forEach((sourceEntityRef, uniqueKey) => {
if (!newRelationSources.has(uniqueKey)) {
setOfThingsToStitch.add(sourceEntityRef);
}
});
await this.stitcher.stitch(setOfThingsToStitch);
track.markSuccessfulWithChanges(setOfThingsToStitch.size);
} catch (error) {
assertError(error);
track.markFailed(error);
}
});
},
});
}
@@ -194,10 +194,12 @@ describe('DefaultCatalogProcessingOrchestrator', () => {
it('runs all processor validations when asked to', async () => {
const validate = jest.fn(async () => true);
const processor1: Partial<CatalogProcessor> = {
const processor1: CatalogProcessor = {
getProcessorName: () => 'processor1',
validateEntityKind: validate,
};
const processor2: Partial<CatalogProcessor> = {
const processor2: CatalogProcessor = {
getProcessorName: () => 'processor2',
validateEntityKind: validate,
};
@@ -14,6 +14,7 @@
* limitations under the License.
*/
import { Span, trace } from '@opentelemetry/api';
import {
Entity,
EntityPolicy,
@@ -55,6 +56,13 @@ import {
} from './util';
import { CatalogRulesEnforcer } from '../ingestion/CatalogRules';
import { ProcessorCacheManager } from './ProcessorCacheManager';
import {
addEntityAttributes,
TRACER_ID,
withActiveSpan,
} from '../util/opentelemetry';
const tracer = trace.getTracer(TRACER_ID);
type Context = {
entityRef: string;
@@ -64,6 +72,18 @@ type Context = {
cache: ProcessorCacheManager;
};
function addProcessorAttributes(
span: Span,
stage: string,
processor: CatalogProcessor,
) {
span.setAttribute('backstage.catalog.processor.stage', stage);
span.setAttribute(
'backstage.catalog.processor.name',
processor.getProcessorName(),
);
}
/** @public */
export class DefaultCatalogProcessingOrchestrator
implements CatalogProcessingOrchestrator
@@ -179,54 +199,71 @@ export class DefaultCatalogProcessingOrchestrator
entity: Entity,
context: Context,
): Promise<Entity> {
let res = entity;
return await withActiveSpan(tracer, 'ProcessingStage', async stageSpan => {
addEntityAttributes(stageSpan, entity);
stageSpan.setAttribute('backstage.catalog.processor.stage', 'preProcess');
let res = entity;
for (const processor of this.options.processors) {
if (processor.preProcessEntity) {
try {
res = await processor.preProcessEntity(
res,
context.location,
context.collector.forProcessor(processor),
context.originLocation,
context.cache.forProcessor(processor),
);
} catch (e) {
throw new InputError(
`Processor ${processor.constructor.name} threw an error while preprocessing`,
e,
);
for (const processor of this.options.processors) {
if (processor.preProcessEntity) {
let innerRes = res;
res = await withActiveSpan(tracer, 'ProcessingStep', async span => {
addEntityAttributes(span, entity);
addProcessorAttributes(span, 'preProcessEntity', processor);
try {
innerRes = await processor.preProcessEntity!(
innerRes,
context.location,
context.collector.forProcessor(processor),
context.originLocation,
context.cache.forProcessor(processor),
);
} catch (e) {
throw new InputError(
`Processor ${processor.constructor.name} threw an error while preprocessing`,
e,
);
}
return innerRes;
});
}
}
}
return res;
return res;
});
}
/**
* Enforce entity policies making sure that entities conform to a general schema
*/
private async runPolicyStep(entity: Entity): Promise<Entity> {
let policyEnforcedEntity: Entity | undefined;
try {
policyEnforcedEntity = await this.options.policy.enforce(entity);
} catch (e) {
throw new InputError(
`Policy check failed for ${stringifyEntityRef(entity)}`,
e,
return await withActiveSpan(tracer, 'ProcessingStage', async stageSpan => {
addEntityAttributes(stageSpan, entity);
stageSpan.setAttribute(
'backstage.catalog.processor.stage',
'enforcePolicy',
);
}
let policyEnforcedEntity: Entity | undefined;
if (!policyEnforcedEntity) {
throw new Error(
`Policy unexpectedly returned no data for ${stringifyEntityRef(
entity,
)}`,
);
}
try {
policyEnforcedEntity = await this.options.policy.enforce(entity);
} catch (e) {
throw new InputError(
`Policy check failed for ${stringifyEntityRef(entity)}`,
e,
);
}
return policyEnforcedEntity;
if (!policyEnforcedEntity) {
throw new Error(
`Policy unexpectedly returned no data for ${stringifyEntityRef(
entity,
)}`,
);
}
return policyEnforcedEntity;
});
}
/**
@@ -236,50 +273,62 @@ export class DefaultCatalogProcessingOrchestrator
entity: Entity,
context: Context,
): Promise<void> {
// Double check that none of the previous steps tried to change something
// related to the entity ref, which would break downstream
if (stringifyEntityRef(entity) !== context.entityRef) {
throw new ConflictError(
'Fatal: The entity kind, namespace, or name changed during processing',
);
}
return await withActiveSpan(tracer, 'ProcessingStage', async stageSpan => {
addEntityAttributes(stageSpan, entity);
stageSpan.setAttribute('backstage.catalog.processor.stage', 'validate');
// Double check that none of the previous steps tried to change something
// related to the entity ref, which would break downstream
if (stringifyEntityRef(entity) !== context.entityRef) {
throw new ConflictError(
'Fatal: The entity kind, namespace, or name changed during processing',
);
}
// Validate that the end result is a valid Entity at all
try {
validateEntity(entity);
} catch (e) {
throw new ConflictError(
`Entity envelope for ${context.entityRef} failed validation after preprocessing`,
e,
);
}
// Validate that the end result is a valid Entity at all
try {
validateEntity(entity);
} catch (e) {
throw new ConflictError(
`Entity envelope for ${context.entityRef} failed validation after preprocessing`,
e,
);
}
let valid = false;
let valid = false;
for (const processor of this.options.processors) {
if (processor.validateEntityKind) {
try {
const thisValid = await processor.validateEntityKind(entity);
if (thisValid) {
valid = true;
if (this.options.legacySingleProcessorValidation) {
break;
for (const processor of this.options.processors) {
if (processor.validateEntityKind) {
try {
const thisValid = await withActiveSpan(
tracer,
'ProcessingStep',
async span => {
addEntityAttributes(span, entity);
addProcessorAttributes(span, 'validateEntityKind', processor);
return await processor.validateEntityKind!(entity);
},
);
if (thisValid) {
valid = true;
if (this.options.legacySingleProcessorValidation) {
break;
}
}
} catch (e) {
throw new InputError(
`Processor ${processor.constructor.name} threw an error while validating the entity ${context.entityRef}`,
e,
);
}
} catch (e) {
throw new InputError(
`Processor ${processor.constructor.name} threw an error while validating the entity ${context.entityRef}`,
e,
);
}
}
}
if (!valid) {
throw new InputError(
`No processor recognized the entity ${context.entityRef} as valid, possibly caused by a foreign kind or apiVersion`,
);
}
if (!valid) {
throw new InputError(
`No processor recognized the entity ${context.entityRef} as valid, possibly caused by a foreign kind or apiVersion`,
);
}
});
}
/**
@@ -289,65 +338,81 @@ export class DefaultCatalogProcessingOrchestrator
entity: LocationEntity,
context: Context,
): Promise<void> {
const { type = context.location.type, presence = 'required' } = entity.spec;
const targets = new Array<string>();
if (entity.spec.target) {
targets.push(entity.spec.target);
}
if (entity.spec.targets) {
targets.push(...entity.spec.targets);
}
for (const maybeRelativeTarget of targets) {
if (type === 'file' && maybeRelativeTarget.endsWith(path.sep)) {
context.collector.generic()(
processingResult.inputError(
context.location,
`LocationEntityProcessor cannot handle ${type} type location with target ${context.location.target} that ends with a path separator`,
),
);
continue;
}
const target = toAbsoluteUrl(
this.options.integrations,
context.location,
type,
maybeRelativeTarget,
return await withActiveSpan(tracer, 'ProcessingStage', async stageSpan => {
addEntityAttributes(stageSpan, entity);
stageSpan.setAttribute(
'backstage.catalog.processor.stage',
'readLocation',
);
const { type = context.location.type, presence = 'required' } =
entity.spec;
const targets = new Array<string>();
if (entity.spec.target) {
targets.push(entity.spec.target);
}
if (entity.spec.targets) {
targets.push(...entity.spec.targets);
}
let didRead = false;
for (const processor of this.options.processors) {
if (processor.readLocation) {
try {
const read = await processor.readLocation(
{
type,
target,
presence,
},
presence === 'optional',
context.collector.forProcessor(processor),
this.options.parser,
context.cache.forProcessor(processor, target),
);
if (read) {
didRead = true;
break;
for (const maybeRelativeTarget of targets) {
if (type === 'file' && maybeRelativeTarget.endsWith(path.sep)) {
context.collector.generic()(
processingResult.inputError(
context.location,
`LocationEntityProcessor cannot handle ${type} type location with target ${context.location.target} that ends with a path separator`,
),
);
continue;
}
const target = toAbsoluteUrl(
this.options.integrations,
context.location,
type,
maybeRelativeTarget,
);
let didRead = false;
for (const processor of this.options.processors) {
if (processor.readLocation) {
try {
const read = await withActiveSpan(
tracer,
'ProcessingStep',
async span => {
addEntityAttributes(span, entity);
addProcessorAttributes(span, 'readLocation', processor);
return await processor.readLocation!(
{
type,
target,
presence,
},
presence === 'optional',
context.collector.forProcessor(processor),
this.options.parser,
context.cache.forProcessor(processor, target),
);
},
);
if (read) {
didRead = true;
break;
}
} catch (e) {
throw new InputError(
`Processor ${processor.constructor.name} threw an error while reading ${type}:${target}`,
e,
);
}
} catch (e) {
throw new InputError(
`Processor ${processor.constructor.name} threw an error while reading ${type}:${target}`,
e,
);
}
}
if (!didRead) {
throw new InputError(
`No processor was able to handle reading of ${type}:${target}`,
);
}
}
if (!didRead) {
throw new InputError(
`No processor was able to handle reading of ${type}:${target}`,
);
}
}
});
}
/**
@@ -357,26 +422,39 @@ export class DefaultCatalogProcessingOrchestrator
entity: Entity,
context: Context,
): Promise<Entity> {
let res = entity;
return await withActiveSpan(tracer, 'ProcessingStage', async stageSpan => {
addEntityAttributes(stageSpan, entity);
stageSpan.setAttribute(
'backstage.catalog.processor.stage',
'postProcessEntity',
);
let res = entity;
for (const processor of this.options.processors) {
if (processor.postProcessEntity) {
try {
res = await processor.postProcessEntity(
res,
context.location,
context.collector.forProcessor(processor),
context.cache.forProcessor(processor),
);
} catch (e) {
throw new InputError(
`Processor ${processor.constructor.name} threw an error while postprocessing`,
e,
);
for (const processor of this.options.processors) {
if (processor.postProcessEntity) {
let innerRes = res;
res = await withActiveSpan(tracer, 'ProcessingStep', async span => {
addEntityAttributes(span, entity);
addProcessorAttributes(span, 'postProcessEntity', processor);
try {
innerRes = await processor.postProcessEntity!(
innerRes,
context.location,
context.collector.forProcessor(processor),
context.cache.forProcessor(processor),
);
} catch (e) {
throw new InputError(
`Processor ${processor.constructor.name} threw an error while postprocessing`,
e,
);
}
return innerRes;
});
}
}
}
return res;
return res;
});
}
}
@@ -0,0 +1,101 @@
/*
* Copyright 2023 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 { Span, SpanOptions, SpanStatusCode, Tracer } from '@opentelemetry/api';
import { Entity } from '@backstage/catalog-model';
export const TRACER_ID = 'backstage-plugin-catalog-backend';
function setAttributeIfDefined(span: Span, attribute: string, value?: string) {
if (value !== null && value !== undefined) {
span.setAttribute(attribute, value);
}
}
export function addEntityAttributes(span: Span, entity: Entity) {
setAttributeIfDefined(span, 'backstage.entity.apiVersion', entity.apiVersion);
setAttributeIfDefined(span, 'backstage.entity.kind', entity.kind);
setAttributeIfDefined(
span,
'backstage.entity.metadata.namespace',
entity.metadata?.namespace,
);
setAttributeIfDefined(
span,
'backstage.entity.metadata.name',
entity.metadata?.name,
);
}
// Adapted from https://github.com/open-telemetry/opentelemetry-js/blob/359fbcc40a859057a02b14e84599eac399b8dba7/api/src/trace/SugaredTracer.ts
// While waiting for something like https://github.com/open-telemetry/opentelemetry-js/pull/3317 to land upstream
const onException = (e: Error, span: Span) => {
span.recordException(e);
span.setStatus({
code: SpanStatusCode.ERROR,
});
};
function isPromiseLike<T, S>(obj: PromiseLike<T> | S): obj is PromiseLike<T> {
return (
!!obj &&
(typeof obj === 'object' || typeof obj === 'function') &&
'then' in obj &&
typeof obj.then === 'function'
);
}
function handleFn<F extends (span: Span) => ReturnType<F>>(
span: Span,
fn: F,
): ReturnType<F> {
try {
const ret = fn(span);
// if fn is an async function attach a recordException and spanEnd callback to the promise
if (isPromiseLike(ret)) {
ret.then(
() => {
span.end();
},
e => {
onException(e, span);
span.end();
},
);
} else {
span.end();
}
return ret;
} catch (e) {
onException(e, span);
span.end();
throw e;
}
}
export function withActiveSpan<F extends (span: Span) => ReturnType<F>>(
tracer: Tracer,
name: string,
fn: F,
spanOptions: SpanOptions = {},
): ReturnType<F> {
return tracer.startActiveSpan(name, spanOptions, (span: Span) => {
return handleFn(span, fn);
});
}