diff --git a/.changeset/late-turtles-listen.md b/.changeset/late-turtles-listen.md new file mode 100644 index 0000000000..117737a064 --- /dev/null +++ b/.changeset/late-turtles-listen.md @@ -0,0 +1,51 @@ +--- +'@backstage/backend-common': patch +--- + +Implemented KubernetesContainerRunner: a ContainerRunner implementation that leverages Jobs on a kubernetes cluster + +```ts +const kubeConfig = new KubeConfig(); +kubeConfig.loadFromDefault(); + +const options: KubernetesContainerRunnerOptions = { + kubeConfig, + // namespace where Jobs will be created + namespace: 'default', + // Jobs name will be prefixed with this name + name: 'my-runner', + // An existing Kubernetes volume that will be used + // as base for mounts + mountBase: { + volumeName: 'workdir', + // Every mount must start with the base path + // see example below + basePath: '/workdir', + }, + // Define a Pod template for the Jobs. It has to include + // a volume definition named as the mountBase volumeName + podTemplate: { + spec: { + containers: [], + volumes: [ + { + name: 'workdir', + persistentVolumeClaim: { + claimName: 'workdir-claim', + }, + }, + ], + }, + }, +}; +const containerRunner = new KubernetesContainerRunner(options); + +const runOptions: RunContainerOptions = { + imageName: 'golang:1.17', + args: ['echo', 'hello world'], + mountDirs: { + '/workdir/app': '/app', + }, +}; +containerRunner.runContainer(runOptions); +``` diff --git a/.github/workflows/verify_kubernetes.yml b/.github/workflows/verify_kubernetes.yml new file mode 100644 index 0000000000..7c692c4457 --- /dev/null +++ b/.github/workflows/verify_kubernetes.yml @@ -0,0 +1,40 @@ +name: E2E Test Kubernetes +on: + pull_request: + paths: + - 'yarn.lock' + - '.github/workflows/verify_kubernetes.yml' + - 'packages/backend-common/src/**' + +jobs: + verify: + runs-on: ubuntu-latest + + strategy: + matrix: + node-version: [14.x, 16.x] + + env: + CI: true + KUBERNETES_TESTS: true + + steps: + - uses: actions/checkout@v3 + + - name: use node.js ${{ matrix.node-version }} + uses: actions/setup-node@v3 + with: + node-version: ${{ matrix.node-version }} + registry-url: https://registry.npmjs.org/ # Needed for auth + + - name: yarn install + uses: backstage/actions/yarn-install@v0.5.3 + with: + cache-prefix: ${{ runner.os }}-v${{ matrix.node-version }} + + - name: bootstrap kind + uses: helm/kind-action@v1.3.0 + + - name: kubernetes test + working-directory: packages/backend-common + run: yarn test:kubernetes diff --git a/packages/backend-common/package.json b/packages/backend-common/package.json index 5f61eb85a6..622f6ae997 100644 --- a/packages/backend-common/package.json +++ b/packages/backend-common/package.json @@ -30,7 +30,8 @@ "prepack": "backstage-cli package prepack", "postpack": "backstage-cli package postpack", "clean": "backstage-cli package clean", - "start": "backstage-cli package start" + "start": "backstage-cli package start", + "test:kubernetes": "backstage-cli package test -t KubernetesContainerRunner --no-watch" }, "dependencies": { "@backstage/cli-common": "^0.1.9", @@ -41,6 +42,7 @@ "@backstage/types": "^1.0.0", "@google-cloud/storage": "^6.0.0", "@keyv/redis": "^2.2.3", + "@kubernetes/client-node": "^0.17.0", "@manypkg/get-packages": "^1.1.3", "@octokit/rest": "^19.0.3", "@types/cors": "^2.8.6", @@ -74,9 +76,11 @@ "node-abort-controller": "^3.0.1", "node-fetch": "^2.6.7", "raw-body": "^2.4.1", + "request": "^2.88.2", "selfsigned": "^2.0.0", "stoppable": "^1.1.0", "tar": "^6.1.2", + "uuid": "^8.3.2", "winston": "^3.2.1", "yauzl": "^2.10.0", "yn": "^4.0.0" diff --git a/packages/backend-common/src/util/KubernetesContainerRunner.test.ts b/packages/backend-common/src/util/KubernetesContainerRunner.test.ts new file mode 100644 index 0000000000..e218abf960 --- /dev/null +++ b/packages/backend-common/src/util/KubernetesContainerRunner.test.ts @@ -0,0 +1,290 @@ +/* + * Copyright 2020 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 { + CoreV1Api, + KubeConfig, + RbacAuthorizationV1Api, +} from '@kubernetes/client-node'; +import { + KubernetesContainerRunner, + KubernetesContainerRunnerOptions, +} from './KubernetesContainerRunner'; +import { RunContainerOptions } from './ContainerRunner'; +import { PassThrough } from 'stream'; + +// This ensures E2E tests that require a Kubernetes cluster are only run +// where KUBERNETES_TESTS environment variable is true +const describeIfKubernetes = Boolean(process.env.KUBERNETES_TESTS) + ? describe + : describe.skip; + +jest.setTimeout(10 * 1000); + +describeIfKubernetes('KubernetesContainerRunner', () => { + const kubeConfig = new KubeConfig(); + kubeConfig.loadFromDefault(); + const name = 'kube-runner'; + + it('should throw error when no namespace is configured', () => { + const testConfig = new KubeConfig(); + testConfig.loadFromDefault(); + testConfig.addContext({ + name: 'test', + cluster: kubeConfig.getCurrentCluster()!.name, + user: kubeConfig.getCurrentUser()!.name, + }); + testConfig.setCurrentContext('test'); + const test = () => + new KubernetesContainerRunner({ + kubeConfig: testConfig, + name, + }); + expect(test).toThrow( + /Cannot read current namespace from Kubernetes cluster/, + ); + }); + + it('should throw error when mountBase is provided and podTemplate is invalid', () => { + const error = /A Pod template containing the volume .+ is required/; + const options: KubernetesContainerRunnerOptions = { + kubeConfig, + name, + namespace: 'default', + mountBase: { + basePath: '/workdir', + volumeName: 'workdir', + }, + }; + const test = () => { + return new KubernetesContainerRunner(options); + }; + expect(test).toThrow(error); + options.podTemplate = {}; + expect(test).toThrow(error); + options.podTemplate.spec = { containers: [] }; + expect(test).toThrow(error); + options.podTemplate.spec.volumes = []; + expect(test).toThrow(error); + }); + + it('should not run the container when the mounts are not subdirectories of the basePath', async () => { + const options: KubernetesContainerRunnerOptions = { + kubeConfig, + name, + namespace: 'default', + mountBase: { + basePath: '/workdir', + volumeName: 'workdir', + }, + podTemplate: { + spec: { + containers: [], + volumes: [ + { + name: 'workdir', + }, + ], + }, + }, + }; + const containerRunner = new KubernetesContainerRunner(options); + const logStream = new PassThrough(); + const runOptions: RunContainerOptions = { + imageName: 'golang:1.17', + args: ['echo', 'hello world'], + logStream, + mountDirs: { + '/notWorkdir/app': '/app', + }, + }; + await expect(containerRunner.runContainer(runOptions)).rejects.toThrow( + `Mounted '/notWorkdir/app' dir should be subdirectories of '${ + options!.mountBase!.basePath + }'`, + ); + }); + + it('should succeed when the container command returns a 0 exit code', async () => { + const options: KubernetesContainerRunnerOptions = { + kubeConfig, + name, + namespace: 'default', + }; + const containerRunner = new KubernetesContainerRunner(options); + const logStream = new PassThrough(); + const chunks: any[] = []; + logStream.on('data', chunk => chunks.push(Buffer.from(chunk))); + const runOptions: RunContainerOptions = { + imageName: 'alpine', + args: ['echo', 'hello world'], + logStream, + }; + await containerRunner.runContainer(runOptions); + const result = Buffer.concat(chunks).toString('utf8'); + expect(result).toBe('hello world\n'); + }); + + it('should fail when container run time exceeds the timeout', async () => { + const options: KubernetesContainerRunnerOptions = { + kubeConfig, + name, + namespace: 'default', + timeoutMs: 5000, + }; + const containerRunner = new KubernetesContainerRunner(options); + const runOptions: RunContainerOptions = { + imageName: 'alpine', + args: ['sleep', '10'], + }; + await expect(containerRunner.runContainer(runOptions)).rejects.toThrow( + `Failed to complete in ${options.timeoutMs} ms`, + ); + }); + + it('should fail when container command returns a non 0 exit code', async () => { + const options: KubernetesContainerRunnerOptions = { + kubeConfig, + name, + namespace: 'default', + }; + const containerRunner = new KubernetesContainerRunner(options); + const runOptions: RunContainerOptions = { + imageName: 'alpine', + args: ['fake'], + }; + await expect(containerRunner.runContainer(runOptions)).rejects.toThrow( + `Container execution failed`, + ); + }); + + it('should fail when job creation fails', async () => { + const options: KubernetesContainerRunnerOptions = { + kubeConfig, + name, + namespace: 'fake', + }; + const containerRunner = new KubernetesContainerRunner(options); + const runOptions: RunContainerOptions = { + imageName: 'golang:1.17', + args: ['echo', 'hello world'], + }; + await expect(containerRunner.runContainer(runOptions)).rejects.toThrow( + 'Kubernetes Job creation failed with the following error message: namespaces "fake" not found', + ); + }); + + describe('with namespace test', () => { + let api: CoreV1Api; + let authApi: RbacAuthorizationV1Api; + + beforeAll(async () => { + api = kubeConfig.makeApiClient(CoreV1Api); + authApi = kubeConfig.makeApiClient(RbacAuthorizationV1Api); + await api.createNamespace({ + metadata: { + name: 'test', + }, + }); + }); + afterAll(async () => { + await api.deleteNamespace('test'); + }); + + it('should fail when watch fails', async () => { + const testConfig = await givenAServiceAccountThatCannotWatchPods( + api, + authApi, + kubeConfig, + ); + const options: KubernetesContainerRunnerOptions = { + kubeConfig: testConfig, + name, + namespace: 'test', + }; + const containerRunner = new KubernetesContainerRunner(options); + const runOptions: RunContainerOptions = { + imageName: 'golang:1.17', + args: ['echo', 'hello world'], + }; + await expect(containerRunner.runContainer(runOptions)).rejects.toThrow( + 'Kubernetes watch request failed with the following error message: Error: Forbidden', + ); + }); + }); +}); + +async function givenAServiceAccountThatCannotWatchPods( + api: CoreV1Api, + authApi: RbacAuthorizationV1Api, + kubeConfig: KubeConfig, +) { + await Promise.all([ + api.createNamespacedServiceAccount('test', { + metadata: { + name: 'test', + }, + }), + authApi.createNamespacedRole('test', { + metadata: { + name: 'test', + }, + rules: [ + { + apiGroups: ['batch'], + verbs: ['create'], + resources: ['jobs'], + }, + ], + }), + authApi.createNamespacedRoleBinding('test', { + metadata: { + name: 'test', + }, + subjects: [ + { + kind: 'ServiceAccount', + name: 'test', + }, + ], + roleRef: { + apiGroup: 'rbac.authorization.k8s.io', + kind: 'Role', + name: 'test', + }, + }), + ]); + const token = ( + await api.createNamespacedServiceAccountToken('test', 'test', { + spec: { + audiences: [], + }, + }) + ).body.status?.token; + const testConfig = new KubeConfig(); + testConfig.loadFromDefault(); + testConfig.addUser({ + name: 'test', + token, + }); + testConfig.addContext({ + name: 'test', + cluster: kubeConfig.getCurrentCluster()!.name, + user: 'test', + }); + testConfig.setCurrentContext('test'); + return testConfig; +} diff --git a/packages/backend-common/src/util/KubernetesContainerRunner.ts b/packages/backend-common/src/util/KubernetesContainerRunner.ts new file mode 100644 index 0000000000..af00351fd1 --- /dev/null +++ b/packages/backend-common/src/util/KubernetesContainerRunner.ts @@ -0,0 +1,379 @@ +/* + * Copyright 2020 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 { PassThrough, Writable } from 'stream'; +import { ContainerRunner, RunContainerOptions } from './ContainerRunner'; +import { + KubeConfig, + BatchV1Api, + V1Job, + V1EnvVar, + Log, + HttpError, + V1Status, + V1VolumeMount, + V1PodTemplateSpec, + V1Pod, + Watch, +} from '@kubernetes/client-node'; +import { v4 as uuid } from 'uuid'; +import { Request } from 'request'; + +/** + * An existing Kubernetes volume that will be used as base for mounts. + * + * Every mount must start with the 'basePath'. + * + * @public + */ +export type MountBaseOptions = { + volumeName: string; + basePath: string; +}; + +/** + * Options to create a {@link KubernetesContainerRunner} + * + * Kubernetes Jobs will be created on the provided 'namespace' + * and their names will be prefixed with the provided 'name'. + * + * 'podTemplate' defines a Pod template for the Jobs. It has to include + * a volume definition named as the {@link MountBaseOptions} 'volumeName'. + * + * @public + */ +export type KubernetesContainerRunnerOptions = { + kubeConfig: KubeConfig; + name: string; + namespace?: string; + mountBase?: MountBaseOptions; + podTemplate?: V1PodTemplateSpec; + timeoutMs?: number; +}; + +/** + * A {@link ContainerRunner} for Kubernetes. + * + * Runs containers leveraging Jobs on a Kubernetes cluster + * + * @public + */ +export class KubernetesContainerRunner implements ContainerRunner { + private readonly kubeConfig: KubeConfig; + private readonly batchV1Api: BatchV1Api; + private readonly log: Log; + private readonly name: string; + private readonly namespace: string; + private readonly mountBase?: MountBaseOptions; + private readonly podTemplate?: V1PodTemplateSpec; + private readonly timeoutMs: number; + private readonly containerName = 'executor'; + + getNamespace(kubeConfig: KubeConfig, namespace?: string): string { + let _namespace = namespace; + if (!_namespace) { + _namespace = kubeConfig.getContextObject( + kubeConfig.currentContext, + )?.namespace; + } + if (!_namespace) { + throw new Error('Cannot read current namespace from Kubernetes cluster'); + } + return _namespace; + } + + validateMountBase( + mountBase: MountBaseOptions, + podTemplate?: V1PodTemplateSpec, + ): MountBaseOptions { + if ( + !podTemplate?.spec?.volumes?.filter(v => v.name === mountBase.volumeName) + .length + ) { + throw new Error( + `A Pod template containing the volume ${mountBase.volumeName} is required`, + ); + } + if (!mountBase.basePath.endsWith('/')) { + mountBase.basePath += '/'; + } + return mountBase; + } + + constructor(options: KubernetesContainerRunnerOptions) { + const { kubeConfig, name, namespace, mountBase, podTemplate, timeoutMs } = + options; + this.kubeConfig = kubeConfig; + this.batchV1Api = kubeConfig.makeApiClient(BatchV1Api); + this.log = new Log(kubeConfig); + this.name = name; + this.namespace = this.getNamespace(kubeConfig, namespace); + if (mountBase) { + this.mountBase = this.validateMountBase(mountBase, podTemplate); + } + this.podTemplate = podTemplate; + this.timeoutMs = timeoutMs || 120 * 1000; + } + + async runContainer(options: RunContainerOptions) { + const { + imageName, + command, + args, + logStream = new PassThrough(), + mountDirs = {}, + workingDir, + envVars = {}, + } = options; + + const commandArr = typeof command === 'string' ? [command] : command; + + const volumeMounts: V1VolumeMount[] = []; + for (const [hostDir, containerDir] of Object.entries(mountDirs)) { + if (!this.mountBase) { + throw new Error( + 'A volumeName and a basePath must be configured to bind mount directories', + ); + } + if (!hostDir.startsWith(this.mountBase.basePath)) { + throw new Error( + `Mounted '${hostDir}' dir should be subdirectories of '${this.mountBase.basePath}'`, + ); + } + volumeMounts.push({ + name: this.mountBase.volumeName, + mountPath: containerDir, + subPath: hostDir.slice(this.mountBase.basePath.length), + }); + } + + const env = []; + for (const [key, value] of Object.entries(envVars)) { + env.push({ + name: key, + value: value, + } as V1EnvVar); + } + + const taskId = uuid(); + + // TODO find a way to merge recursively + const mergedPodTemplate: V1PodTemplateSpec = { + metadata: { + ...{ + labels: { + task: taskId, + }, + }, + ...this.podTemplate?.metadata, + }, + spec: { + ...{ + containers: [ + { + name: this.containerName, + image: imageName, + command: commandArr, + args: args, + env: env, + workingDir: workingDir, + volumeMounts: volumeMounts, + }, + ], + restartPolicy: 'Never', + }, + ...this.podTemplate?.spec, + }, + }; + + const jobSpec: V1Job = { + metadata: { + generateName: `${this.name}-`, + }, + spec: { + backoffLimit: 0, + ttlSecondsAfterFinished: 60, + template: mergedPodTemplate, + }, + }; + + await this.runJob(jobSpec, taskId, logStream); + } + + handleError(err: any, errorCallback: (reason: any) => void) { + if (err.code !== 'ECONNRESET' && err.message !== 'aborted') { + errorCallback( + handleKubernetesError( + 'Kubernetes watch request failed with the following error message:', + err, + ), + ); + } + } + + watchPod( + taskId: string, + callback: (pod: V1Pod) => void, + errorCallback: (reason: any) => void, + ): Promise { + const watch = new Watch(this.kubeConfig); + const labelSelector = `task=${taskId}`; + return watch.watch( + `/api/v1/namespaces/${this.namespace}/pods`, + { + labelSelector, + }, + (_, pod) => { + callback(pod); + }, + err => { + if (err) { + this.handleError(err, errorCallback); + } + }, + ); + } + + tailLogs( + taskId: string, + logStream: Writable, + ): { promise: Promise; close: () => Promise } { + let log: Promise; + let req: Promise; + const watchPromise = new Promise((_, reject) => { + req = this.watchPod( + taskId, + pod => { + if ( + log === undefined && + (pod.status?.phase === 'Running' || + pod.status?.phase === 'Succeeded' || + pod.status?.phase === 'Failed') + ) { + log = this.log.log( + this.namespace, + pod.metadata?.name!, + this.containerName, + logStream, + { follow: true }, + ); + } + }, + reject, + ); + }); + const logPromise = new Promise((resolve, _) => { + if (!logStream.writableFinished) { + logStream.on('finish', () => { + resolve(); + }); + } else { + resolve(); + } + }); + const close = async () => { + if (req) { + (await req).abort(); + } + if (log) { + (await log).abort(); + } + }; + return { promise: Promise.race([watchPromise, logPromise]), close }; + } + + waitPod(taskId: string): { + promise: Promise; + close: () => Promise; + } { + let req: Promise; + const promise = new Promise(async (resolve, reject) => { + req = this.watchPod( + taskId, + pod => { + if (pod.status?.phase === 'Succeeded') { + resolve(); + } + if (pod.status?.phase === 'Failed') { + reject(new Error('Container execution failed')); + } + }, + reject, + ); + }); + const close = async () => { + if (req) { + (await req).abort(); + } + }; + return { promise, close }; + } + + async createJob(jobSpec: V1Job): Promise { + return this.batchV1Api + .createNamespacedJob(this.namespace, jobSpec) + .catch(err => { + throw handleKubernetesError( + 'Kubernetes Job creation failed with the following error message:', + err, + ); + }); + } + + async runJob( + jobSpec: V1Job, + taskId: string, + logStream: Writable, + ): Promise { + let timeout: NodeJS.Timeout; + const timeoutPromise = new Promise((_, reject) => { + timeout = setTimeout( + reject, + this.timeoutMs, + new Error(`Failed to complete in ${this.timeoutMs} ms`), + ); + }); + + const { promise: waitPromise, close: waitClose } = this.waitPod(taskId); + const { promise: tailPromise, close: tailClose } = this.tailLogs( + taskId, + logStream, + ); + + const taskPromise = Promise.all([ + waitPromise, + tailPromise, + this.createJob(jobSpec), + ]).finally(() => { + clearTimeout(timeout); + }); + + return Promise.race([timeoutPromise, taskPromise]) + .finally(() => { + return waitClose(); + }) + .finally(() => { + return tailClose(); + }); + } +} + +function handleKubernetesError(message: string, err: Error): Error { + if (err instanceof HttpError) { + return new Error(`${message} ${(err.body as V1Status).message}`); + } + return new Error(`${message} ${err}`); +} diff --git a/yarn.lock b/yarn.lock index f0a567bfe3..7a6081a3cc 100644 --- a/yarn.lock +++ b/yarn.lock @@ -2867,6 +2867,7 @@ __metadata: "@backstage/types": ^1.0.0 "@google-cloud/storage": ^6.0.0 "@keyv/redis": ^2.2.3 + "@kubernetes/client-node": ^0.17.0 "@manypkg/get-packages": ^1.1.3 "@octokit/rest": ^19.0.3 "@types/archiver": ^5.1.0 @@ -2921,10 +2922,12 @@ __metadata: node-fetch: ^2.6.7 raw-body: ^2.4.1 recursive-readdir: ^2.2.2 + request: ^2.88.2 selfsigned: ^2.0.0 stoppable: ^1.1.0 supertest: ^6.1.3 tar: ^6.1.2 + uuid: ^8.3.2 winston: ^3.2.1 yauzl: ^2.10.0 yn: ^4.0.0 @@ -35750,7 +35753,7 @@ __metadata: languageName: node linkType: hard -"request@npm:^2.88.0": +"request@npm:^2.88.0, request@npm:^2.88.2": version: 2.88.2 resolution: "request@npm:2.88.2" dependencies: