Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 21 additions & 0 deletions plugins/kubernetes-ingestor/src/extensions.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
import { createExtensionPoint } from '@backstage/backend-plugin-api';
import { KubernetesResourceFilter } from './types';

/**
* Extension point for narrowing what the ingestor puts in the catalog.
*
* The configuration options cover the two common cases — a namespace list and an
* annotation on the object itself — but neither reaches resources that are spread
* across every namespace and deployed from charts the adopter does not own. A filter
* is code, so it can decide on any part of the resource: name, labels, owner
* references, whatever the installation needs.
*/
export interface KubernetesIngestorExtensionPoint {
/** Registers one or more filters. All of them must pass for a resource to be ingested. */
addResourceFilter(...filters: KubernetesResourceFilter[]): void;
}

export const kubernetesIngestorExtensionPoint =
createExtensionPoint<KubernetesIngestorExtensionPoint>({
id: 'kubernetes-ingestor',
});
7 changes: 7 additions & 0 deletions plugins/kubernetes-ingestor/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,3 +2,10 @@ export { catalogModuleKubernetesIngestor as default } from './module';
export { KubernetesEntityProvider, XRDTemplateEntityProvider } from './providers';
export type { DeltaEvent } from './providers';
export type { KubernetesResourceFetcher, KubernetesResourceFetcherOptions } from './types';
export { kubernetesIngestorExtensionPoint } from './extensions';
export type { KubernetesIngestorExtensionPoint } from './extensions';
export type {
KubernetesResourceFilter,
KubernetesResourceFilterContext,
KubernetesResourceFilterInput,
} from './types';
10 changes: 10 additions & 0 deletions plugins/kubernetes-ingestor/src/module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@ import {
import { EventParams, eventsServiceRef } from '@backstage/plugin-events-node';
import { KubernetesEntityProvider, RGDTemplateEntityProvider, XRDTemplateEntityProvider } from './providers';
import { DefaultKubernetesResourceFetcher } from './services';
import { kubernetesIngestorExtensionPoint } from './extensions';
import { KubernetesResourceFilter } from './types';

interface DeltaEventPayload {
action: 'upsert' | 'delete';
Expand Down Expand Up @@ -57,6 +59,13 @@ export const catalogModuleKubernetesIngestor = createBackendModule({
pluginId: 'catalog',
moduleId: 'kubernetes-ingestor',
register(reg) {
const resourceFilters: KubernetesResourceFilter[] = [];
reg.registerExtensionPoint(kubernetesIngestorExtensionPoint, {
addResourceFilter(...filters: KubernetesResourceFilter[]) {
resourceFilters.push(...filters);
},
});

reg.registerInit({
deps: {
catalog: catalogProcessingExtensionPoint,
Expand Down Expand Up @@ -115,6 +124,7 @@ export const catalogModuleKubernetesIngestor = createBackendModule({
resourceFetcher,
urlReader,
cache,
resourceFilters,
);

const xrdTemplateEntityProvider = new XRDTemplateEntityProvider(
Expand Down
157 changes: 157 additions & 0 deletions plugins/kubernetes-ingestor/src/providers/EntityProvider.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2386,6 +2386,126 @@ describe('KubernetesEntityProvider', () => {
expect(deltaCalls[0][0].removed).toEqual([]);
});

it('should remove entities when a delta upsert becomes ineligible', async () => {
let isEligible = true;
const resourceFilter = jest.fn().mockImplementation(() => isEligible);
const provider = new KubernetesEntityProvider(
{ run: jest.fn() } as any,
mockLogger,
mockConfig,
mockResourceFetcher as any,
undefined,
undefined,
[resourceFilter],
);

const mockConnection = {
applyMutation: jest.fn().mockResolvedValue(undefined),
};

await provider.connect(mockConnection as any);
(provider as any).fullSyncCompleted = true;

mockResourceFetcher.proxyKubernetesRequest.mockResolvedValue({
apiVersion: 'apps/v1',
kind: 'Deployment',
metadata: {
name: 'filtered-deployment',
namespace: 'default',
},
spec: {},
});

await provider.deltaUpdate({
action: 'upsert',
apiVersion: 'apps/v1',
kind: 'Deployment',
name: 'filtered-deployment',
namespace: 'default',
clusterName: 'test-cluster',
});

isEligible = false;

await provider.deltaUpdate({
action: 'upsert',
apiVersion: 'apps/v1',
kind: 'Deployment',
name: 'filtered-deployment',
namespace: 'default',
clusterName: 'test-cluster',
});

expect(resourceFilter).toHaveBeenLastCalledWith(
expect.objectContaining({
metadata: expect.objectContaining({ name: 'filtered-deployment' }),
}),
{ clusterName: 'test-cluster' },
);
const deltaCalls = mockConnection.applyMutation.mock.calls.filter(
(call: any[]) => call[0].type === 'delta',
);
expect(deltaCalls).toHaveLength(2);
expect(deltaCalls[0][0].added.length).toBeGreaterThan(0);
expect(deltaCalls[0][0].removed).toEqual([]);
expect(deltaCalls[1][0].added).toEqual([]);
expect(deltaCalls[1][0].removed.length).toBeGreaterThan(0);
expect(
deltaCalls[1][0].removed.some(
(entry: any) => entry.entity.kind === 'System',
),
).toBe(false);
});

it('should not run filters for a disabled Crossplane delta upsert', async () => {
const resourceFilter = jest.fn().mockReturnValue(true);
const config = new ConfigReader({
kubernetesIngestor: {
components: { enabled: true },
crossplane: { enabled: false },
kro: { enabled: false },
annotationPrefix: 'terasky.backstage.io',
},
});
const provider = new KubernetesEntityProvider(
{ run: jest.fn() } as any,
mockLogger,
config,
mockResourceFetcher as any,
undefined,
undefined,
[resourceFilter],
);

const mockConnection = {
applyMutation: jest.fn().mockResolvedValue(undefined),
};

await provider.connect(mockConnection as any);
(provider as any).fullSyncCompleted = true;
mockResourceFetcher.proxyKubernetesRequest.mockResolvedValueOnce({
apiVersion: 'example.org/v1',
kind: 'Example',
metadata: { name: 'crossplane-resource', namespace: 'default' },
spec: { crossplane: {} },
});

await provider.deltaUpdate({
action: 'upsert',
apiVersion: 'example.org/v1',
kind: 'Example',
name: 'crossplane-resource',
namespace: 'default',
clusterName: 'test-cluster',
});

expect(resourceFilter).not.toHaveBeenCalled();
const deltaCalls = mockConnection.applyMutation.mock.calls.filter(
(call: any[]) => call[0].type === 'delta',
);
expect(deltaCalls).toHaveLength(0);
});

it('should perform delta delete for a regular K8s resource', async () => {
const provider = new KubernetesEntityProvider(
{ run: jest.fn() } as any,
Expand Down Expand Up @@ -2419,6 +2539,43 @@ describe('KubernetesEntityProvider', () => {
expect(deltaCalls[0][0].removed.length).toBeGreaterThan(0);
});

it('should not apply resource filters to delta deletes', async () => {
const resourceFilter = jest.fn().mockReturnValue(false);
const provider = new KubernetesEntityProvider(
{ run: jest.fn() } as any,
mockLogger,
mockConfig,
mockResourceFetcher as any,
undefined,
undefined,
[resourceFilter],
);

const mockConnection = {
applyMutation: jest.fn().mockResolvedValue(undefined),
};

await provider.connect(mockConnection as any);
(provider as any).fullSyncCompleted = true;

await provider.deltaUpdate({
action: 'delete',
apiVersion: 'apps/v1',
kind: 'Deployment',
name: 'filtered-deployment',
namespace: 'default',
clusterName: 'test-cluster',
});

expect(resourceFilter).not.toHaveBeenCalled();
const deltaCalls = mockConnection.applyMutation.mock.calls.filter(
(call: any[]) => call[0].type === 'delta',
);
expect(deltaCalls).toHaveLength(1);
expect(deltaCalls[0][0].added).toEqual([]);
expect(deltaCalls[0][0].removed.length).toBeGreaterThan(0);
});

it('should handle resource fetch failure gracefully on upsert', async () => {
const provider = new KubernetesEntityProvider(
{ run: jest.fn() } as any,
Expand Down
53 changes: 52 additions & 1 deletion plugins/kubernetes-ingestor/src/providers/EntityProvider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,14 @@ import { Entity, parseEntityRef } from '@backstage/catalog-model';
import { Config } from '@backstage/config';
import { CacheService, LoggerService, SchedulerServiceTaskRunner, UrlReaderService } from '@backstage/backend-plugin-api';
import { DefaultKubernetesResourceFetcher, ApiDefinitionFetcher } from '../services';
import { KubernetesDataProvider, WorkloadType } from './KubernetesDataProvider';
import { KubernetesResourceFilter } from '../types';
import {
createBuiltInResourceEligibilityChecker,
getDisabledResourceType,
KubernetesDataProvider,
passesResourceFilters,
WorkloadType,
} from './KubernetesDataProvider';
import { Logger } from 'winston';
import { CRDDataProvider } from './CRDDataProvider';
import { XRDDataProvider } from './XRDDataProvider';
Expand Down Expand Up @@ -2438,6 +2445,7 @@ export class KubernetesEntityProvider implements EntityProvider {
private readonly resourceFetcher: DefaultKubernetesResourceFetcher,
urlReader?: UrlReaderService,
private readonly cache?: CacheService,
private readonly resourceFilters: KubernetesResourceFilter[] = [],
) {
this.orphanGraceRuns = config.getOptionalNumber(
'kubernetesIngestor.orphanProtection.gracePeriodRuns',
Expand Down Expand Up @@ -2702,6 +2710,7 @@ export class KubernetesEntityProvider implements EntityProvider {
this.resourceFetcher,
this.config,
this.logger,
this.resourceFilters,
);

let compositeKindLookup: { [key: string]: any } = {};
Expand Down Expand Up @@ -2982,6 +2991,48 @@ export class KubernetesEntityProvider implements EntityProvider {
return;
}
resource.clusterName = clusterName;

const passesBuiltInEligibility =
createBuiltInResourceEligibilityChecker(this.config);
if (!passesBuiltInEligibility(resource)) {
this.logger.debug(
`Skipping delta upsert for ${kind}/${name}: resource is not eligible for ingestion`,
);
return;
}

const disabledResourceType = getDisabledResourceType(
resource,
isCrossplaneEnabled,
isKROEnabled,
);
if (disabledResourceType) {
this.logger.debug(
`Skipping delta upsert for ${kind}/${name}: ${disabledResourceType} ingestion is disabled`,
);
return;
}

if (
!passesResourceFilters(
resource,
{ clusterName },
this.resourceFilters,
)
) {
const { entities } = await this.classifyAndTranslateResource(
resource, isCrossplaneEnabled, isKROEnabled,
this.cachedCompositeKindLookup, this.cachedRgdLookup, this.cachedCrdMapping,
);
const removable = entities.filter(entity => entity.kind !== 'System');
if (removable.length > 0) {
await this.applyDeltaMutation([], removable);
}
this.logger.info(
`Delta upsert removed ${removable.length} entities for ineligible ${kind}/${name} from cluster ${clusterName}`,
);
return;
}
} else {
// For deletes without explicit entityNames, construct a synthetic resource
// enriched with enough metadata for classifyAndTranslateResource to route
Expand Down
Loading