Merge branch 'backstage:master' into ldap-entity-search-fix

This commit is contained in:
John
2024-11-23 09:56:12 +11:00
committed by GitHub
17 changed files with 484 additions and 33 deletions
+5
View File
@@ -0,0 +1,5 @@
---
'@backstage/backend-defaults': patch
---
Export `PluginTokenHandler` and `pluginTokenHandlerDecoratorServiceRef` to allow for custom decoration of the plugin token handler without having to re-implement the entire handler.
+198
View File
@@ -0,0 +1,198 @@
{
"mode": "pre",
"tag": "next",
"initialVersions": {
"example-app": "0.2.103",
"@backstage/app-defaults": "1.5.13",
"example-app-next": "0.0.17",
"app-next-example-plugin": "0.0.17",
"example-backend": "0.0.32",
"@backstage/backend-app-api": "1.0.2",
"@backstage/backend-defaults": "0.5.3",
"@backstage/backend-dev-utils": "0.1.5",
"@backstage/backend-dynamic-feature-service": "0.5.0",
"example-backend-legacy": "0.2.104",
"@backstage/backend-openapi-utils": "0.3.0",
"@backstage/backend-plugin-api": "1.0.2",
"@backstage/backend-test-utils": "1.1.0",
"@backstage/catalog-client": "1.8.0",
"@backstage/catalog-model": "1.7.1",
"@backstage/cli": "0.29.0",
"@backstage/cli-common": "0.1.15",
"@backstage/cli-node": "0.2.10",
"@backstage/codemods": "0.1.52",
"@backstage/config": "1.3.0",
"@backstage/config-loader": "1.9.2",
"@backstage/core-app-api": "1.15.2",
"@backstage/core-compat-api": "0.3.2",
"@backstage/core-components": "0.16.0",
"@backstage/core-plugin-api": "1.10.1",
"@backstage/create-app": "0.5.22",
"@backstage/dev-utils": "1.1.3",
"e2e-test": "0.2.22",
"@backstage/e2e-test-utils": "0.1.1",
"@backstage/errors": "1.2.5",
"@backstage/eslint-plugin": "0.1.10",
"@backstage/frontend-app-api": "0.10.1",
"@backstage/frontend-defaults": "0.1.2",
"@internal/frontend": "0.0.3",
"@backstage/frontend-plugin-api": "0.9.1",
"@backstage/frontend-test-utils": "0.2.2",
"@backstage/integration": "1.15.2",
"@backstage/integration-aws-node": "0.1.13",
"@backstage/integration-react": "1.2.1",
"@internal/opaque": "0.0.1",
"@backstage/release-manifests": "0.0.11",
"@backstage/repo-tools": "0.11.0",
"@internal/scaffolder": "0.0.3",
"@techdocs/cli": "1.8.22",
"techdocs-cli-embedded-app": "0.2.102",
"@backstage/test-utils": "1.7.1",
"@backstage/theme": "0.6.1",
"@backstage/types": "1.2.0",
"@backstage/version-bridge": "1.0.10",
"yarn-plugin-backstage": "0.0.3",
"@backstage/plugin-api-docs": "0.12.0",
"@backstage/plugin-api-docs-module-protoc-gen-doc": "0.1.8",
"@backstage/plugin-app": "0.1.2",
"@backstage/plugin-app-backend": "0.4.0",
"@backstage/plugin-app-node": "0.1.27",
"@backstage/plugin-app-visualizer": "0.1.12",
"@backstage/plugin-auth-backend": "0.24.0",
"@backstage/plugin-auth-backend-module-atlassian-provider": "0.3.2",
"@backstage/plugin-auth-backend-module-auth0-provider": "0.1.2",
"@backstage/plugin-auth-backend-module-aws-alb-provider": "0.3.0",
"@backstage/plugin-auth-backend-module-azure-easyauth-provider": "0.2.2",
"@backstage/plugin-auth-backend-module-bitbucket-provider": "0.2.2",
"@backstage/plugin-auth-backend-module-bitbucket-server-provider": "0.1.2",
"@backstage/plugin-auth-backend-module-cloudflare-access-provider": "0.3.2",
"@backstage/plugin-auth-backend-module-gcp-iap-provider": "0.3.2",
"@backstage/plugin-auth-backend-module-github-provider": "0.2.2",
"@backstage/plugin-auth-backend-module-gitlab-provider": "0.2.2",
"@backstage/plugin-auth-backend-module-google-provider": "0.2.2",
"@backstage/plugin-auth-backend-module-guest-provider": "0.2.2",
"@backstage/plugin-auth-backend-module-microsoft-provider": "0.2.2",
"@backstage/plugin-auth-backend-module-oauth2-provider": "0.3.2",
"@backstage/plugin-auth-backend-module-oauth2-proxy-provider": "0.2.2",
"@backstage/plugin-auth-backend-module-oidc-provider": "0.3.2",
"@backstage/plugin-auth-backend-module-okta-provider": "0.1.2",
"@backstage/plugin-auth-backend-module-onelogin-provider": "0.2.2",
"@backstage/plugin-auth-backend-module-pinniped-provider": "0.2.2",
"@backstage/plugin-auth-backend-module-vmware-cloud-provider": "0.4.1",
"@backstage/plugin-auth-node": "0.5.4",
"@backstage/plugin-auth-react": "0.1.8",
"@backstage/plugin-bitbucket-cloud-common": "0.2.25",
"@backstage/plugin-catalog": "1.25.0",
"@backstage/plugin-catalog-backend": "1.28.0",
"@backstage/plugin-catalog-backend-module-aws": "0.4.5",
"@backstage/plugin-catalog-backend-module-azure": "0.2.4",
"@backstage/plugin-catalog-backend-module-backstage-openapi": "0.4.2",
"@backstage/plugin-catalog-backend-module-bitbucket-cloud": "0.4.2",
"@backstage/plugin-catalog-backend-module-bitbucket-server": "0.2.4",
"@backstage/plugin-catalog-backend-module-gcp": "0.3.2",
"@backstage/plugin-catalog-backend-module-gerrit": "0.2.4",
"@backstage/plugin-catalog-backend-module-github": "0.7.7",
"@backstage/plugin-catalog-backend-module-github-org": "0.3.4",
"@backstage/plugin-catalog-backend-module-gitlab": "0.5.0",
"@backstage/plugin-catalog-backend-module-gitlab-org": "0.2.3",
"@backstage/plugin-catalog-backend-module-incremental-ingestion": "0.6.0",
"@backstage/plugin-catalog-backend-module-ldap": "0.10.0",
"@backstage/plugin-catalog-backend-module-logs": "0.1.4",
"@backstage/plugin-catalog-backend-module-msgraph": "0.6.4",
"@backstage/plugin-catalog-backend-module-openapi": "0.2.4",
"@backstage/plugin-catalog-backend-module-puppetdb": "0.2.4",
"@backstage/plugin-catalog-backend-module-scaffolder-entity-model": "0.2.2",
"@backstage/plugin-catalog-backend-module-unprocessed": "0.5.2",
"@backstage/plugin-catalog-common": "1.1.1",
"@backstage/plugin-catalog-graph": "0.4.12",
"@backstage/plugin-catalog-import": "0.12.6",
"@backstage/plugin-catalog-node": "1.14.0",
"@backstage/plugin-catalog-react": "1.14.1",
"@backstage/plugin-catalog-unprocessed-entities": "0.2.10",
"@backstage/plugin-catalog-unprocessed-entities-common": "0.0.5",
"@backstage/plugin-config-schema": "0.1.61",
"@backstage/plugin-devtools": "0.1.20",
"@backstage/plugin-devtools-backend": "0.4.2",
"@backstage/plugin-devtools-common": "0.1.13",
"@backstage/plugin-events-backend": "0.3.16",
"@backstage/plugin-events-backend-module-aws-sqs": "0.4.5",
"@backstage/plugin-events-backend-module-azure": "0.2.14",
"@backstage/plugin-events-backend-module-bitbucket-cloud": "0.2.14",
"@backstage/plugin-events-backend-module-gerrit": "0.2.14",
"@backstage/plugin-events-backend-module-github": "0.2.14",
"@backstage/plugin-events-backend-module-gitlab": "0.2.14",
"@backstage/plugin-events-backend-test-utils": "0.1.38",
"@backstage/plugin-events-node": "0.4.5",
"@internal/plugin-todo-list": "1.0.33",
"@internal/plugin-todo-list-backend": "1.0.33",
"@internal/plugin-todo-list-common": "1.0.22",
"@backstage/plugin-home": "0.8.1",
"@backstage/plugin-home-react": "0.1.19",
"@backstage/plugin-kubernetes": "0.12.0",
"@backstage/plugin-kubernetes-backend": "0.19.0",
"@backstage/plugin-kubernetes-cluster": "0.0.18",
"@backstage/plugin-kubernetes-common": "0.9.0",
"@backstage/plugin-kubernetes-node": "0.2.0",
"@backstage/plugin-kubernetes-react": "0.5.0",
"@backstage/plugin-notifications": "0.4.0",
"@backstage/plugin-notifications-backend": "0.4.3",
"@backstage/plugin-notifications-backend-module-email": "0.3.3",
"@backstage/plugin-notifications-common": "0.0.6",
"@backstage/plugin-notifications-node": "0.2.9",
"@backstage/plugin-org": "0.6.32",
"@backstage/plugin-org-react": "0.1.31",
"@backstage/plugin-permission-backend": "0.5.51",
"@backstage/plugin-permission-backend-module-allow-all-policy": "0.2.2",
"@backstage/plugin-permission-common": "0.8.2",
"@backstage/plugin-permission-node": "0.8.5",
"@backstage/plugin-permission-react": "0.4.28",
"@backstage/plugin-proxy-backend": "0.5.8",
"@backstage/plugin-scaffolder": "1.27.0",
"@backstage/plugin-scaffolder-backend": "1.27.0",
"@backstage/plugin-scaffolder-backend-module-azure": "0.2.2",
"@backstage/plugin-scaffolder-backend-module-bitbucket": "0.3.2",
"@backstage/plugin-scaffolder-backend-module-bitbucket-cloud": "0.2.2",
"@backstage/plugin-scaffolder-backend-module-bitbucket-server": "0.2.2",
"@backstage/plugin-scaffolder-backend-module-confluence-to-markdown": "0.3.2",
"@backstage/plugin-scaffolder-backend-module-cookiecutter": "0.3.3",
"@backstage/plugin-scaffolder-backend-module-gcp": "0.2.2",
"@backstage/plugin-scaffolder-backend-module-gerrit": "0.2.2",
"@backstage/plugin-scaffolder-backend-module-gitea": "0.2.2",
"@backstage/plugin-scaffolder-backend-module-github": "0.5.2",
"@backstage/plugin-scaffolder-backend-module-gitlab": "0.6.1",
"@backstage/plugin-scaffolder-backend-module-notifications": "0.1.3",
"@backstage/plugin-scaffolder-backend-module-rails": "0.5.2",
"@backstage/plugin-scaffolder-backend-module-sentry": "0.2.2",
"@backstage/plugin-scaffolder-backend-module-yeoman": "0.4.3",
"@backstage/plugin-scaffolder-common": "1.5.7",
"@backstage/plugin-scaffolder-node": "0.6.0",
"@backstage/plugin-scaffolder-node-test-utils": "0.1.15",
"@backstage/plugin-scaffolder-react": "1.14.0",
"@backstage/plugin-search": "1.4.19",
"@backstage/plugin-search-backend": "1.7.0",
"@backstage/plugin-search-backend-module-catalog": "0.2.5",
"@backstage/plugin-search-backend-module-elasticsearch": "1.6.2",
"@backstage/plugin-search-backend-module-explore": "0.2.5",
"@backstage/plugin-search-backend-module-pg": "0.5.38",
"@backstage/plugin-search-backend-module-stack-overflow-collator": "0.3.3",
"@backstage/plugin-search-backend-module-techdocs": "0.3.2",
"@backstage/plugin-search-backend-node": "1.3.5",
"@backstage/plugin-search-common": "1.2.15",
"@backstage/plugin-search-react": "1.8.2",
"@backstage/plugin-signals": "0.0.12",
"@backstage/plugin-signals-backend": "0.2.3",
"@backstage/plugin-signals-node": "0.1.14",
"@backstage/plugin-signals-react": "0.0.7",
"@backstage/plugin-techdocs": "1.11.1",
"@backstage/plugin-techdocs-addons-test-utils": "1.0.41",
"@backstage/plugin-techdocs-backend": "1.11.2",
"@backstage/plugin-techdocs-common": "0.1.0",
"@backstage/plugin-techdocs-module-addons-contrib": "1.1.17",
"@backstage/plugin-techdocs-node": "1.12.13",
"@backstage/plugin-techdocs-react": "1.2.10",
"@backstage/plugin-user-settings": "0.8.15",
"@backstage/plugin-user-settings-backend": "0.2.27",
"@backstage/plugin-user-settings-common": "0.0.1"
},
"changesets": []
}
+5
View File
@@ -0,0 +1,5 @@
---
'@backstage/plugin-catalog-backend-module-incremental-ingestion': patch
---
Wire up the events together in the new backend system
+28
View File
@@ -412,3 +412,31 @@ Each entry has one or more of the following fields:
# Also supports the shorthand form:
# action: create, read
```
## Adding custom or logic for validation and issuing of tokens
The `pluginTokenHandlerDecoratorServiceRef` can be used to decorate the existing token handler without having to re-implement the entire `AuthService` implementation.
This is particularly useful when you want to add additional logic to the handler, such as logging or metrics or custom token validation.
The `PluginTokenHandler` interface has two methods:
- `issueToken`: This method is used to issue a token for a plugin. It takes in the `pluginId` and `targetPluginId` as arguments, and an optional `limitedUserToken` object which can be used to issue a token on behalf of another user. The method returns a promise that resolves to an object containing the issued token.
- `verifyToken`: This method is used to verify a token. It takes in the token as an argument and returns a promise that resolves to an object containing the subject of the token and an optional limited user token.
```ts
import {
PluginTokenHandler,
pluginTokenHandlerDecoratorServiceRef,
} from '@backstage/backend-defaults/auth';
import { createServiceFactory } from '@backstage/backend-plugin-api';
const decoratedPluginTokenHandler = createServiceFactory({
service: pluginTokenHandlerDecoratorServiceRef,
deps: {},
async factory() {
return (defaultImplementation: PluginTokenHandler) =>
new CustomTokenHandler(defaultImplementation);
},
});
```
@@ -5,6 +5,7 @@
```ts
import { AuthService } from '@backstage/backend-plugin-api';
import { ServiceFactory } from '@backstage/backend-plugin-api';
import { ServiceRef } from '@backstage/backend-plugin-api';
// @public
export const authServiceFactory: ServiceFactory<
@@ -13,5 +14,35 @@ export const authServiceFactory: ServiceFactory<
'singleton'
>;
// @public
export interface PluginTokenHandler {
// (undocumented)
issueToken(options: {
pluginId: string;
targetPluginId: string;
onBehalfOf?: {
limitedUserToken: string;
expiresAt: Date;
};
}): Promise<{
token: string;
}>;
// (undocumented)
verifyToken(token: string): Promise<
| {
subject: string;
limitedUserToken?: string;
}
| undefined
>;
}
// @public
export const pluginTokenHandlerDecoratorServiceRef: ServiceRef<
(defaultImplementation: PluginTokenHandler) => PluginTokenHandler,
'plugin',
'singleton'
>;
// (No @packageDocumentation comment for this package)
```
@@ -169,7 +169,10 @@ export class DefaultAuthService implements AuthService {
return this.pluginTokenHandler.issueToken({
pluginId: this.pluginId,
targetPluginId,
onBehalfOf,
onBehalfOf: {
limitedUserToken: onBehalfOf.token,
expiresAt: onBehalfOf.expiresAt,
},
});
}
default:
@@ -19,12 +19,17 @@ import {
mockServices,
registerMswTestHooks,
} from '@backstage/backend-test-utils';
import { authServiceFactory } from './authServiceFactory';
import {
authServiceFactory,
pluginTokenHandlerDecoratorServiceRef,
} from './authServiceFactory';
import { base64url, decodeJwt } from 'jose';
import { discoveryServiceFactory } from '../discovery';
import { rest } from 'msw';
import { setupServer } from 'msw/node';
import { toInternalBackstageCredentials } from './helpers';
import { PluginTokenHandler } from './plugin/PluginTokenHandler';
import { createServiceFactory } from '@backstage/backend-plugin-api';
const server = setupServer();
@@ -407,4 +412,40 @@ describe('authServiceFactory', () => {
principal: { subject: 'unlimited-static-subject' },
});
});
describe('decorate PluginTokenHandler', () => {
it('should allow custom logic to be injected into the plugin token handler', async () => {
const customLogic = jest.fn();
const customPluginTokenHandler = createServiceFactory({
service: pluginTokenHandlerDecoratorServiceRef,
deps: {},
async factory() {
return (defaultImplementation: PluginTokenHandler) =>
new (class CustomHandler implements PluginTokenHandler {
verifyToken(
token: string,
): Promise<
{ subject: string; limitedUserToken?: string } | undefined
> {
customLogic(token);
return defaultImplementation.verifyToken(token);
}
issueToken(options: {
pluginId: string;
targetPluginId: string;
limitedUserToken?: { token: string; expiresAt: Date };
}): Promise<{ token: string }> {
return defaultImplementation.issueToken(options);
}
})();
},
});
const tester = ServiceFactoryTester.from(authServiceFactory, {
dependencies: [...mockDeps, customPluginTokenHandler],
});
const searchAuth = await tester.getSubject('search');
searchAuth.authenticate('unlimited-static-token');
expect(customLogic).toHaveBeenCalledWith('unlimited-static-token');
});
});
});
@@ -17,13 +17,35 @@
import {
coreServices,
createServiceFactory,
createServiceRef,
} from '@backstage/backend-plugin-api';
import { DefaultAuthService } from './DefaultAuthService';
import { ExternalTokenHandler } from './external/ExternalTokenHandler';
import { PluginTokenHandler } from './plugin/PluginTokenHandler';
import {
DefaultPluginTokenHandler,
PluginTokenHandler,
} from './plugin/PluginTokenHandler';
import { createPluginKeySource } from './plugin/keys/createPluginKeySource';
import { UserTokenHandler } from './user/UserTokenHandler';
/**
* @public
* This service is used to decorate the default plugin token handler with custom logic.
*/
export const pluginTokenHandlerDecoratorServiceRef = createServiceRef<
(defaultImplementation: PluginTokenHandler) => PluginTokenHandler
>({
id: 'core.auth.pluginTokenHandlerDecorator',
defaultFactory: async service =>
createServiceFactory({
service,
deps: {},
factory: async () => {
return impl => impl;
},
}),
});
/**
* Handles token authentication and credentials management.
*
@@ -41,8 +63,16 @@ export const authServiceFactory = createServiceFactory({
discovery: coreServices.discovery,
plugin: coreServices.pluginMetadata,
database: coreServices.database,
pluginTokenHandlerDecorator: pluginTokenHandlerDecoratorServiceRef,
},
async factory({ config, discovery, plugin, logger, database }) {
async factory({
config,
discovery,
plugin,
logger,
database,
pluginTokenHandlerDecorator,
}) {
const disableDefaultAuthPolicy =
config.getOptionalBoolean(
'backend.auth.dangerouslyDisableDefaultAuthPolicy',
@@ -61,13 +91,15 @@ export const authServiceFactory = createServiceFactory({
discovery,
});
const pluginTokens = PluginTokenHandler.create({
ownPluginId: plugin.getId(),
logger,
keySource,
keyDuration,
discovery,
});
const pluginTokens = pluginTokenHandlerDecorator(
DefaultPluginTokenHandler.create({
ownPluginId: plugin.getId(),
logger,
keySource,
keyDuration,
discovery,
}),
);
const externalTokens = ExternalTokenHandler.create({
ownPluginId: plugin.getId(),
@@ -14,4 +14,9 @@
* limitations under the License.
*/
export { authServiceFactory } from './authServiceFactory';
export {
authServiceFactory,
pluginTokenHandlerDecoratorServiceRef,
} from './authServiceFactory';
export type { PluginTokenHandler } from './plugin/PluginTokenHandler';
@@ -15,7 +15,7 @@
*/
import { mockServices } from '@backstage/backend-test-utils';
import { PluginTokenHandler } from './PluginTokenHandler';
import { DefaultPluginTokenHandler } from './PluginTokenHandler';
import { decodeJwt } from 'jose';
describe('PluginTokenHandler', () => {
@@ -42,7 +42,7 @@ describe('PluginTokenHandler', () => {
});
const getKeyMock = jest.fn(async () => mockPrivateKey);
const handler = PluginTokenHandler.create({
const handler = DefaultPluginTokenHandler.create({
discovery: mockServices.discovery(),
keyDuration: { seconds: 10 },
logger: mockServices.logger.mock(),
@@ -16,13 +16,12 @@
import { DiscoveryService, LoggerService } from '@backstage/backend-plugin-api';
import { decodeJwt, importJWK, SignJWT, decodeProtectedHeader } from 'jose';
import { AuthenticationError } from '@backstage/errors';
import { assertError, AuthenticationError } from '@backstage/errors';
import { jwtVerify } from 'jose';
import { tokenTypes } from '@backstage/plugin-auth-node';
import { JwksClient } from '../JwksClient';
import { HumanDuration, durationToMilliseconds } from '@backstage/types';
import { PluginKeySource } from './keys/types';
import fetch from 'node-fetch';
const SECONDS_IN_MS = 1000;
@@ -45,7 +44,22 @@ type Options = {
algorithm?: string;
};
export class PluginTokenHandler {
/**
* @public
* Issues and verifies {@link https://backstage.iceio/docs/auth/service-to-service-auth | service-to-service tokens}.
*/
export interface PluginTokenHandler {
verifyToken(
token: string,
): Promise<{ subject: string; limitedUserToken?: string } | undefined>;
issueToken(options: {
pluginId: string;
targetPluginId: string;
onBehalfOf?: { limitedUserToken: string; expiresAt: Date };
}): Promise<{ token: string }>;
}
export class DefaultPluginTokenHandler implements PluginTokenHandler {
private jwksMap = new Map<string, JwksClient>();
// Tracking state for isTargetPluginSupported
@@ -53,7 +67,7 @@ export class PluginTokenHandler {
private targetPluginInflightChecks = new Map<string, Promise<boolean>>();
static create(options: Options) {
return new PluginTokenHandler(
return new DefaultPluginTokenHandler(
options.logger,
options.ownPluginId,
options.keySource,
@@ -115,7 +129,7 @@ export class PluginTokenHandler {
async issueToken(options: {
pluginId: string;
targetPluginId: string;
onBehalfOf?: { token: string; expiresAt: Date };
onBehalfOf?: { limitedUserToken: string; expiresAt: Date };
}): Promise<{ token: string }> {
const { pluginId, targetPluginId, onBehalfOf } = options;
const key = await this.keySource.getPrivateSigningKey();
@@ -131,7 +145,7 @@ export class PluginTokenHandler {
)
: ourExp;
const claims = { sub, aud, iat, exp, obo: onBehalfOf?.token };
const claims = { sub, aud, iat, exp, obo: onBehalfOf?.limitedUserToken };
const token = await new SignJWT(claims)
.setProtectedHeader({
typ: tokenTypes.plugin.typParam,
@@ -181,6 +195,7 @@ export class PluginTokenHandler {
this.supportedTargetPlugins.add(targetPluginId);
return true;
} catch (error) {
assertError(error);
this.logger.error('Unexpected failure for target JWKS check', error);
return false;
} finally {
-6
View File
@@ -16,9 +16,3 @@
// package: the name of the package, e.g. @testing-library/react
// query: the version query to pin the version for, e.g. ^14.0.0
// version: the version to pin to, must be in range of the query, e.g. 14.11.0
// https://github.com/swagger-api/apidom/issues/4539
"@swagger-api/apidom-json-pointer@^1.0.0-alpha.1":
version "1.0.0-alpha.10"
resolved "https://registry.yarnpkg.com/@swagger-api/apidom-json-pointer/-/apidom-json-pointer-1.0.0-alpha.10.tgz#762e972b9200b603ca4311ed34246d6026d5b03e"
integrity sha512-Xo0v4Jxp0ZiAm+OOL2PSLyjiw5OAkCMxI0nN9+vOw1/mfXcC+tdb30QQ9WNtF7O9LExjznfFID/NnDEYqBRDwA==
@@ -68,6 +68,7 @@ describe('WrapperProviders', () => {
client,
scheduler: scheduler as Partial<SchedulerService> as SchedulerService,
applyDatabaseMigrations,
events: mockServices.events.mock(),
});
const wrapped1 = providers.wrap(provider1, {
burstInterval: { seconds: 1 },
@@ -36,6 +36,7 @@ import {
IncrementalEntityProvider,
IncrementalEntityProviderOptions,
} from '../types';
import { EventsService } from '@backstage/plugin-events-node';
/**
* Helps in the creation of the catalog entity providers that wrap the
@@ -53,6 +54,7 @@ export class WrapperProviders {
client: Knex;
scheduler: SchedulerService;
applyDatabaseMigrations?: typeof applyDatabaseMigrations;
events: EventsService;
},
) {}
@@ -130,6 +132,20 @@ export class WrapperProviders {
frequency,
timeout: length,
});
const topics = engine.supportsEventTopics();
if (topics.length > 0) {
logger.info(
`Provider ${provider.getProviderName()} subscribing to events for topics: ${topics.join(
',',
)}`,
);
await this.options.events.subscribe({
topics,
id: `catalog-backend-module-incremental-ingestion:${provider.getProviderName()}`,
onEvent: evt => engine.onEvent(evt),
});
}
} catch (error) {
logger.warn(
`Failed to initialize incremental ingestion provider ${provider.getProviderName()}, ${stringifyError(
@@ -25,6 +25,7 @@ import {
IncrementalEntityProviderOptions,
} from '@backstage/plugin-catalog-backend-module-incremental-ingestion';
import { WrapperProviders } from './WrapperProviders';
import { eventsServiceRef } from '@backstage/plugin-events-node';
/**
* @public
@@ -106,6 +107,7 @@ export const catalogModuleIncrementalIngestionEntityProvider =
httpRouter: coreServices.httpRouter,
logger: coreServices.logger,
scheduler: coreServices.scheduler,
events: eventsServiceRef,
},
async init({
catalog,
@@ -114,6 +116,7 @@ export const catalogModuleIncrementalIngestionEntityProvider =
httpRouter,
logger,
scheduler,
events,
}) {
const client = await database.getClient();
@@ -122,6 +125,7 @@ export const catalogModuleIncrementalIngestionEntityProvider =
logger,
client,
scheduler,
events,
});
for (const entry of addedProviders) {
@@ -805,6 +805,18 @@ describe('DefaultProviderDatabase', () => {
metadata: { namespace: 'ns', name: 'n4' },
};
const entity5: Entity = {
apiVersion: '1',
kind: 'k',
metadata: { namespace: 'ns', name: 'n5' },
};
const entity6: Entity = {
apiVersion: '1',
kind: 'k',
metadata: { namespace: 'ns', name: 'n6' },
};
await insertRefreshStateRow(knex, {
entity_id: 'id1',
entity_ref: stringifyEntityRef(entity1Before),
@@ -844,6 +856,33 @@ describe('DefaultProviderDatabase', () => {
source_key: 'my-provider',
target_entity_ref: stringifyEntityRef(entity4),
});
await insertRefreshStateRow(knex, {
entity_id: 'id5',
entity_ref: stringifyEntityRef(entity5),
last_discovery_at: new Date(),
next_update_at: new Date(),
errors: '[]',
unprocessed_entity: JSON.stringify(entity5),
unprocessed_hash: generateStableHash(entity5),
});
await insertRefRow(knex, {
source_key: 'my-provider',
target_entity_ref: stringifyEntityRef(entity5),
});
await insertRefreshStateRow(knex, {
entity_id: 'id6',
entity_ref: stringifyEntityRef(entity6),
last_discovery_at: new Date(),
next_update_at: new Date(),
errors: '[]',
unprocessed_entity: JSON.stringify(entity6),
unprocessed_hash: generateStableHash(entity6),
location_key: 'old',
});
await insertRefRow(knex, {
source_key: 'my-provider',
target_entity_ref: stringifyEntityRef(entity6),
});
await db.transaction(async tx => {
await db.replaceUnprocessedEntities(tx, {
@@ -856,13 +895,22 @@ describe('DefaultProviderDatabase', () => {
{ entity: entity2After },
// this didn't exist, so will become an add
{ entity: entity3 },
// only the location key changed, so this will become an update
{ entity: entity5, locationKey: 'new' },
// only the location key changed, so this will become an update
{ entity: entity6, locationKey: 'new' },
],
removed: [{ entityRef: stringifyEntityRef(entity4) }],
});
});
const state = await knex<DbRefreshStateRow>('refresh_state')
.select(['entity_ref', 'unprocessed_entity', 'unprocessed_hash'])
.select([
'entity_ref',
'unprocessed_entity',
'unprocessed_hash',
'location_key',
])
.orderBy('entity_ref');
expect(state).toEqual([
@@ -870,16 +918,32 @@ describe('DefaultProviderDatabase', () => {
entity_ref: stringifyEntityRef(entity1After),
unprocessed_entity: JSON.stringify(entity1After),
unprocessed_hash: generateStableHash(entity1After),
location_key: null,
},
{
entity_ref: stringifyEntityRef(entity2After),
unprocessed_entity: JSON.stringify(entity2Before), // didn't change
unprocessed_hash: generateStableHash(entity2After),
location_key: null,
},
{
entity_ref: stringifyEntityRef(entity3),
unprocessed_entity: JSON.stringify(entity3),
unprocessed_hash: generateStableHash(entity3),
location_key: null,
},
// entity4 was deleted here
{
entity_ref: stringifyEntityRef(entity5),
unprocessed_entity: JSON.stringify(entity5),
unprocessed_hash: generateStableHash(entity5),
location_key: 'new', // permitted to change, because it was null before
},
{
entity_ref: stringifyEntityRef(entity6),
unprocessed_entity: JSON.stringify(entity6),
unprocessed_hash: generateStableHash(entity6),
location_key: 'old', // TODO(freben): This should have been updated to "new", but we do not support that yet
},
]);
},
@@ -220,19 +220,28 @@ export class DefaultProviderDatabase implements ProviderDatabase {
for (const chunk of lodash.chunk(options.added, 1000)) {
const entityRefs = chunk.map(e => stringifyEntityRef(e.entity));
const rows = await tx<DbRefreshStateRow>('refresh_state')
.select(['entity_ref', 'unprocessed_hash'])
.select(['entity_ref', 'unprocessed_hash', 'location_key'])
.whereIn('entity_ref', entityRefs);
const oldHashes = new Map(
rows.map(row => [row.entity_ref, row.unprocessed_hash]),
const oldStates = new Map(
rows.map(row => [
row.entity_ref,
{
unprocessed_hash: row.unprocessed_hash,
location_key: row.location_key,
},
]),
);
chunk.forEach((deferred, i) => {
const entityRef = entityRefs[i];
const newHash = generateStableHash(deferred.entity);
const oldHash = oldHashes.get(entityRef);
if (oldHash === undefined) {
const oldState = oldStates.get(entityRef);
if (oldState === undefined) {
toAdd.push({ deferred, hash: newHash });
} else if (newHash !== oldHash) {
} else if (
newHash !== oldState.unprocessed_hash ||
(deferred.locationKey ?? null) !== (oldState.location_key ?? null) // normalize undefined/null
) {
toUpsert.push({ deferred, hash: newHash });
}
});