From e0a6360b80bd9b77093d1ce37f8b64cce314723b Mon Sep 17 00:00:00 2001 From: Eric Peterson Date: Thu, 28 Apr 2022 19:43:04 +0200 Subject: [PATCH] Add stream() to UrlReader.readUrl response Signed-off-by: Eric Peterson --- .changeset/backend-common-slips-away.md | 7 + packages/backend-common/api-report.md | 18 +++ packages/backend-common/package.json | 2 + .../src/reading/AwsS3UrlReader.test.ts | 13 +- .../src/reading/AwsS3UrlReader.ts | 10 +- .../src/reading/AzureUrlReader.ts | 5 +- .../reading/BitbucketCloudUrlReader.test.ts | 25 +++- .../src/reading/BitbucketCloudUrlReader.ts | 6 +- .../src/reading/BitbucketServerUrlReader.ts | 6 +- .../src/reading/BitbucketUrlReader.test.ts | 25 +++- .../src/reading/BitbucketUrlReader.ts | 6 +- .../src/reading/FetchUrlReader.ts | 6 +- .../src/reading/GerritUrlReader.test.ts | 23 +++- .../src/reading/GerritUrlReader.ts | 14 +- .../src/reading/GithubUrlReader.ts | 6 +- .../src/reading/GitlabUrlReader.ts | 6 +- .../src/reading/GoogleGcsUrlReader.ts | 17 ++- .../reading/ReadUrlResponseFactory.test.ts | 125 ++++++++++++++++++ .../src/reading/ReadUrlResponseFactory.ts | 77 +++++++++++ packages/backend-common/src/reading/index.ts | 2 + packages/backend-common/src/reading/types.ts | 18 +++ yarn.lock | 12 ++ 22 files changed, 388 insertions(+), 41 deletions(-) create mode 100644 .changeset/backend-common-slips-away.md create mode 100644 packages/backend-common/src/reading/ReadUrlResponseFactory.test.ts create mode 100644 packages/backend-common/src/reading/ReadUrlResponseFactory.ts diff --git a/.changeset/backend-common-slips-away.md b/.changeset/backend-common-slips-away.md new file mode 100644 index 0000000000..93b77f1111 --- /dev/null +++ b/.changeset/backend-common-slips-away.md @@ -0,0 +1,7 @@ +--- +'@backstage/backend-common': patch +--- + +Added a `stream()` method to complement the `buffer()` method on `ReadUrlResponse`. A `ReadUrlResponseFactory` utility class is now also available, providing a simple, consistent way to provide a valid `ReadUrlResponse`. + +This method, though optional for now, will be required on the responses of `UrlReader.readUrl()` implementations in a future release. diff --git a/packages/backend-common/api-report.md b/packages/backend-common/api-report.md index 961973b144..13c4ed62ce 100644 --- a/packages/backend-common/api-report.md +++ b/packages/backend-common/api-report.md @@ -547,6 +547,24 @@ export type ReadUrlOptions = { // @public export type ReadUrlResponse = { buffer(): Promise; + stream?(): Readable; + etag?: string; +}; + +// @public +export class ReadUrlResponseFactory { + static fromNodeJSReadable( + oldStyleStream: NodeJS.ReadableStream, + options?: ReadUrlResponseFactoryFromStreamOptions, + ): Promise; + static fromReadable( + stream: Readable, + options?: ReadUrlResponseFactoryFromStreamOptions, + ): Promise; +} + +// @public +export type ReadUrlResponseFactoryFromStreamOptions = { etag?: string; }; diff --git a/packages/backend-common/package.json b/packages/backend-common/package.json index d66465a471..b32cd6551d 100644 --- a/packages/backend-common/package.json +++ b/packages/backend-common/package.json @@ -50,6 +50,7 @@ "@types/webpack-env": "^1.15.2", "archiver": "^5.0.2", "aws-sdk": "^2.840.0", + "base64-stream": "^1.0.0", "compression": "^1.7.4", "concat-stream": "^2.0.0", "cors": "^2.8.5", @@ -93,6 +94,7 @@ "@backstage/backend-test-utils": "^0.1.24-next.0", "@backstage/cli": "^0.17.1-next.0", "@types/archiver": "^5.1.0", + "@types/base64-stream": "^1.0.2", "@types/compression": "^1.7.0", "@types/concat-stream": "^2.0.0", "@types/fs-extra": "^9.0.3", diff --git a/packages/backend-common/src/reading/AwsS3UrlReader.test.ts b/packages/backend-common/src/reading/AwsS3UrlReader.test.ts index 01fd27f45a..ea969b43a9 100644 --- a/packages/backend-common/src/reading/AwsS3UrlReader.test.ts +++ b/packages/backend-common/src/reading/AwsS3UrlReader.test.ts @@ -28,6 +28,7 @@ import AWSMock from 'aws-sdk-mock'; import aws from 'aws-sdk'; import path from 'path'; import { NotModifiedError } from '@backstage/errors'; +import getRawBody from 'raw-body'; const treeResponseFactory = DefaultReadTreeResponseFactory.create({ config: new ConfigReader({}), @@ -263,7 +264,7 @@ describe('AwsS3UrlReader', () => { describe('readUrl', () => { let awsS3UrlReader: AwsS3UrlReader; - beforeAll(() => { + beforeEach(() => { AWSMock.setSDKInstance(aws); AWSMock.mock( @@ -295,7 +296,7 @@ describe('AwsS3UrlReader', () => { ); }); - it('returns contents of an object in a bucket', async () => { + it('returns contents of an object in a bucket via buffer', async () => { const response = await awsS3UrlReader.readUrl( 'https://test-bucket.s3.us-east-2.amazonaws.com/awsS3-mock-object.yaml', ); @@ -303,6 +304,14 @@ describe('AwsS3UrlReader', () => { expect(buffer.toString().trim()).toBe('site_name: Test'); }); + it('returns contents of an object in a bucket via stream', async () => { + const response = await awsS3UrlReader.readUrl( + 'https://test-bucket.s3.us-east-2.amazonaws.com/awsS3-mock-object.yaml', + ); + const fromStream = await getRawBody(response.stream!()); + expect(fromStream.toString().trim()).toBe('site_name: Test'); + }); + it('rejects unknown targets', async () => { await expect( awsS3UrlReader.readUrl( diff --git a/packages/backend-common/src/reading/AwsS3UrlReader.ts b/packages/backend-common/src/reading/AwsS3UrlReader.ts index 4dd68bf32c..dbd2c66a62 100644 --- a/packages/backend-common/src/reading/AwsS3UrlReader.ts +++ b/packages/backend-common/src/reading/AwsS3UrlReader.ts @@ -26,7 +26,6 @@ import { SearchResponse, UrlReader, } from './types'; -import getRawBody from 'raw-body'; import { AwsS3Integration, ScmIntegrations, @@ -34,6 +33,7 @@ import { } from '@backstage/integration'; import { ForwardedError, NotModifiedError } from '@backstage/errors'; import { ListObjectsV2Output, ObjectList } from 'aws-sdk/clients/s3'; +import { ReadUrlResponseFactory } from './ReadUrlResponseFactory'; /** * Path style URLs: https://s3.(region).amazonaws.com/(bucket)/(key) @@ -214,13 +214,11 @@ export class AwsS3UrlReader implements UrlReader { const request = this.deps.s3.getObject(params); options?.signal?.addEventListener('abort', () => request.abort()); - const buffer = await getRawBody(request.createReadStream()); const etag = (await request.promise()).ETag; - return { - buffer: async () => buffer, - etag: etag, - }; + return ReadUrlResponseFactory.fromReadable(request.createReadStream(), { + etag, + }); } catch (e) { if (e.statusCode === 304) { throw new NotModifiedError(); diff --git a/packages/backend-common/src/reading/AzureUrlReader.ts b/packages/backend-common/src/reading/AzureUrlReader.ts index d6ceaec258..140496c181 100644 --- a/packages/backend-common/src/reading/AzureUrlReader.ts +++ b/packages/backend-common/src/reading/AzureUrlReader.ts @@ -37,6 +37,7 @@ import { ReadUrlOptions, ReadUrlResponse, } from './types'; +import { ReadUrlResponseFactory } from './ReadUrlResponseFactory'; /** * Implements a {@link UrlReader} for Azure repos. @@ -90,9 +91,7 @@ export class AzureUrlReader implements UrlReader { // for private repos when PAT is not valid, Azure API returns a http status code 203 with sign in page html if (response.ok && response.status !== 203) { - return { - buffer: async () => Buffer.from(await response.arrayBuffer()), - }; + return ReadUrlResponseFactory.fromNodeJSReadable(response.body); } const message = `${url} could not be read as ${builtUrl}, ${response.status} ${response.statusText}`; diff --git a/packages/backend-common/src/reading/BitbucketCloudUrlReader.test.ts b/packages/backend-common/src/reading/BitbucketCloudUrlReader.test.ts index 38841d75ad..12ace2755b 100644 --- a/packages/backend-common/src/reading/BitbucketCloudUrlReader.test.ts +++ b/packages/backend-common/src/reading/BitbucketCloudUrlReader.test.ts @@ -29,6 +29,7 @@ import path from 'path'; import { NotModifiedError } from '@backstage/errors'; import { BitbucketCloudUrlReader } from './BitbucketCloudUrlReader'; import { DefaultReadTreeResponseFactory } from './tree'; +import getRawBody from 'raw-body'; const treeResponseFactory = DefaultReadTreeResponseFactory.create({ config: new ConfigReader({}), @@ -65,7 +66,7 @@ describe('BitbucketCloudUrlReader', () => { setupRequestMockHandlers(worker); describe('readUrl', () => { - it('should be able to readUrl without ETag', async () => { + it('should be able to readUrl via buffer without ETag', async () => { worker.use( rest.get( 'https://api.bitbucket.org/2.0/repositories/backstage-verification/test-template/src/master/template.yaml', @@ -87,6 +88,28 @@ describe('BitbucketCloudUrlReader', () => { expect(buffer.toString()).toBe('foo'); }); + it('should be able to readUrl via stream without ETag', async () => { + worker.use( + rest.get( + 'https://api.bitbucket.org/2.0/repositories/backstage-verification/test-template/src/master/template.yaml', + (req, res, ctx) => { + expect(req.headers.get('If-None-Match')).toBeNull(); + return res( + ctx.status(200), + ctx.body('foo'), + ctx.set('ETag', 'etag-value'), + ); + }, + ), + ); + + const result = await reader.readUrl( + 'https://bitbucket.org/backstage-verification/test-template/src/master/template.yaml', + ); + const fromStream = await getRawBody(result.stream!()); + expect(fromStream.toString()).toBe('foo'); + }); + it('should be able to readUrl with matching ETag', async () => { worker.use( rest.get( diff --git a/packages/backend-common/src/reading/BitbucketCloudUrlReader.ts b/packages/backend-common/src/reading/BitbucketCloudUrlReader.ts index 6a9406e920..7bc4cd9c11 100644 --- a/packages/backend-common/src/reading/BitbucketCloudUrlReader.ts +++ b/packages/backend-common/src/reading/BitbucketCloudUrlReader.ts @@ -39,6 +39,7 @@ import { SearchResponse, UrlReader, } from './types'; +import { ReadUrlResponseFactory } from './ReadUrlResponseFactory'; /** * Implements a {@link UrlReader} for files from Bitbucket Cloud. @@ -112,10 +113,9 @@ export class BitbucketCloudUrlReader implements UrlReader { } if (response.ok) { - return { - buffer: async () => Buffer.from(await response.arrayBuffer()), + return ReadUrlResponseFactory.fromNodeJSReadable(response.body, { etag: response.headers.get('ETag') ?? undefined, - }; + }); } const message = `${url} could not be read as ${bitbucketUrl}, ${response.status} ${response.statusText}`; diff --git a/packages/backend-common/src/reading/BitbucketServerUrlReader.ts b/packages/backend-common/src/reading/BitbucketServerUrlReader.ts index 320b8f2719..b4e1e62327 100644 --- a/packages/backend-common/src/reading/BitbucketServerUrlReader.ts +++ b/packages/backend-common/src/reading/BitbucketServerUrlReader.ts @@ -38,6 +38,7 @@ import { SearchResponse, UrlReader, } from './types'; +import { ReadUrlResponseFactory } from './ReadUrlResponseFactory'; /** * Implements a {@link UrlReader} for files from Bitbucket Server APIs. @@ -103,10 +104,9 @@ export class BitbucketServerUrlReader implements UrlReader { } if (response.ok) { - return { - buffer: async () => Buffer.from(await response.arrayBuffer()), + return ReadUrlResponseFactory.fromNodeJSReadable(response.body, { etag: response.headers.get('ETag') ?? undefined, - }; + }); } const message = `${url} could not be read as ${bitbucketUrl}, ${response.status} ${response.statusText}`; diff --git a/packages/backend-common/src/reading/BitbucketUrlReader.test.ts b/packages/backend-common/src/reading/BitbucketUrlReader.test.ts index 3a2360debb..922480f8e7 100644 --- a/packages/backend-common/src/reading/BitbucketUrlReader.test.ts +++ b/packages/backend-common/src/reading/BitbucketUrlReader.test.ts @@ -30,6 +30,7 @@ import { NotModifiedError } from '@backstage/errors'; import { BitbucketUrlReader } from './BitbucketUrlReader'; import { DefaultReadTreeResponseFactory } from './tree'; import { getVoidLogger } from '../logging'; +import getRawBody from 'raw-body'; const logger = getVoidLogger(); @@ -113,7 +114,7 @@ describe('BitbucketUrlReader', () => { setupRequestMockHandlers(worker); describe('readUrl', () => { - it('should be able to readUrl without ETag', async () => { + it('should be able to readUrl via buffer without ETag', async () => { worker.use( rest.get( 'https://api.bitbucket.org/2.0/repositories/backstage-verification/test-template/src/master/template.yaml', @@ -135,6 +136,28 @@ describe('BitbucketUrlReader', () => { expect(buffer.toString()).toBe('foo'); }); + it('should be able to readUrl via stream without ETag', async () => { + worker.use( + rest.get( + 'https://api.bitbucket.org/2.0/repositories/backstage-verification/test-template/src/master/template.yaml', + (req, res, ctx) => { + expect(req.headers.get('If-None-Match')).toBeNull(); + return res( + ctx.status(200), + ctx.body('foo'), + ctx.set('ETag', 'etag-value'), + ); + }, + ), + ); + + const result = await bitbucketProcessor.readUrl( + 'https://bitbucket.org/backstage-verification/test-template/src/master/template.yaml', + ); + const fromStream = await getRawBody(result.stream!()); + expect(fromStream.toString()).toBe('foo'); + }); + it('should be able to readUrl with matching ETag', async () => { worker.use( rest.get( diff --git a/packages/backend-common/src/reading/BitbucketUrlReader.ts b/packages/backend-common/src/reading/BitbucketUrlReader.ts index 1687a66262..805697bf27 100644 --- a/packages/backend-common/src/reading/BitbucketUrlReader.ts +++ b/packages/backend-common/src/reading/BitbucketUrlReader.ts @@ -40,6 +40,7 @@ import { SearchResponse, UrlReader, } from './types'; +import { ReadUrlResponseFactory } from './ReadUrlResponseFactory'; /** * Implements a {@link UrlReader} for files from Bitbucket v1 and v2 APIs, such @@ -123,10 +124,9 @@ export class BitbucketUrlReader implements UrlReader { } if (response.ok) { - return { - buffer: async () => Buffer.from(await response.arrayBuffer()), + return ReadUrlResponseFactory.fromNodeJSReadable(response.body, { etag: response.headers.get('ETag') ?? undefined, - }; + }); } const message = `${url} could not be read as ${bitbucketUrl}, ${response.status} ${response.statusText}`; diff --git a/packages/backend-common/src/reading/FetchUrlReader.ts b/packages/backend-common/src/reading/FetchUrlReader.ts index cb873ebdf8..b83f08c22f 100644 --- a/packages/backend-common/src/reading/FetchUrlReader.ts +++ b/packages/backend-common/src/reading/FetchUrlReader.ts @@ -25,6 +25,7 @@ import { UrlReader, } from './types'; import path from 'path'; +import { ReadUrlResponseFactory } from './ReadUrlResponseFactory'; /** * A {@link UrlReader} that does a plain fetch of the URL. @@ -103,10 +104,9 @@ export class FetchUrlReader implements UrlReader { } if (response.ok) { - return { - buffer: async () => Buffer.from(await response.arrayBuffer()), + return ReadUrlResponseFactory.fromNodeJSReadable(response.body, { etag: response.headers.get('ETag') ?? undefined, - }; + }); } const message = `could not read ${url}, ${response.status} ${response.statusText}`; diff --git a/packages/backend-common/src/reading/GerritUrlReader.test.ts b/packages/backend-common/src/reading/GerritUrlReader.test.ts index fc27b3dd6a..4e0b97fc58 100644 --- a/packages/backend-common/src/reading/GerritUrlReader.test.ts +++ b/packages/backend-common/src/reading/GerritUrlReader.test.ts @@ -31,6 +31,7 @@ import { getVoidLogger } from '../logging'; import { UrlReaderPredicateTuple } from './types'; import { DefaultReadTreeResponseFactory } from './tree'; import { GerritUrlReader } from './GerritUrlReader'; +import getRawBody from 'raw-body'; const treeResponseFactory = DefaultReadTreeResponseFactory.create({ config: new ConfigReader({}), @@ -128,7 +129,7 @@ describe('GerritUrlReader', () => { describe('readUrl', () => { const responseBuffer = Buffer.from('Apache License'); - it('should be able to read file contents', async () => { + it('should be able to read file contents as buffer', async () => { worker.use( rest.get( 'https://gerrit.com/projects/web%2Fproject/branches/master/files/LICENSE/content', @@ -148,6 +149,26 @@ describe('GerritUrlReader', () => { expect(buffer.toString()).toBe(responseBuffer.toString()); }); + it('should be able to read file contents as stream', async () => { + worker.use( + rest.get( + 'https://gerrit.com/projects/web%2Fproject/branches/master/files/LICENSE/content', + (_, res, ctx) => { + return res( + ctx.status(200), + ctx.body(responseBuffer.toString('base64')), + ); + }, + ), + ); + + const result = await gerritProcessor.readUrl( + 'https://gerrit.com/web/project/+/refs/heads/master/LICENSE', + ); + const fromStream = await getRawBody(result.stream!()); + expect(fromStream.toString()).toBe(responseBuffer.toString()); + }); + it('should raise NotFoundError on 404.', async () => { worker.use( rest.get( diff --git a/packages/backend-common/src/reading/GerritUrlReader.ts b/packages/backend-common/src/reading/GerritUrlReader.ts index f3e8924192..59b33bb2cb 100644 --- a/packages/backend-common/src/reading/GerritUrlReader.ts +++ b/packages/backend-common/src/reading/GerritUrlReader.ts @@ -25,6 +25,7 @@ import { parseGerritJsonResponse, parseGerritGitilesUrl, } from '@backstage/integration'; +import { Base64Decode } from 'base64-stream'; import concatStream from 'concat-stream'; import fs from 'fs-extra'; import fetch, { Response } from 'node-fetch'; @@ -125,9 +126,18 @@ export class GerritUrlReader implements UrlReader { throw new Error(`Unable to read gerrit file ${url}, ${e}`); } if (response.ok) { - const responseBody = await response.text(); + let responseBody: string; return { - buffer: async () => Buffer.from(responseBody, 'base64'), + buffer: async () => { + if (responseBody === undefined) { + responseBody = await response.text(); + } + return Buffer.from(responseBody, 'base64'); + }, + stream: () => { + const readable = new Readable().wrap(response.body); + return readable.pipe(new Base64Decode()); + }, }; } if (response.status === 404) { diff --git a/packages/backend-common/src/reading/GithubUrlReader.ts b/packages/backend-common/src/reading/GithubUrlReader.ts index d86f491827..3806253a6f 100644 --- a/packages/backend-common/src/reading/GithubUrlReader.ts +++ b/packages/backend-common/src/reading/GithubUrlReader.ts @@ -39,6 +39,7 @@ import { ReadUrlOptions, ReadUrlResponse, } from './types'; +import { ReadUrlResponseFactory } from './ReadUrlResponseFactory'; export type GhRepoResponse = RestEndpointMethodTypes['repos']['get']['response']['data']; @@ -127,10 +128,9 @@ export class GithubUrlReader implements UrlReader { } if (response.ok) { - return { - buffer: async () => Buffer.from(await response.arrayBuffer()), + return ReadUrlResponseFactory.fromNodeJSReadable(response.body, { etag: response.headers.get('ETag') ?? undefined, - }; + }); } let message = `${url} could not be read as ${ghUrl}, ${response.status} ${response.statusText}`; diff --git a/packages/backend-common/src/reading/GitlabUrlReader.ts b/packages/backend-common/src/reading/GitlabUrlReader.ts index d0e43fea59..2387dadf5b 100644 --- a/packages/backend-common/src/reading/GitlabUrlReader.ts +++ b/packages/backend-common/src/reading/GitlabUrlReader.ts @@ -38,6 +38,7 @@ import { ReadUrlOptions, } from './types'; import { trimEnd } from 'lodash'; +import { ReadUrlResponseFactory } from './ReadUrlResponseFactory'; /** * Implements a {@link UrlReader} for files on GitLab. @@ -97,10 +98,9 @@ export class GitlabUrlReader implements UrlReader { } if (response.ok) { - return { - buffer: async () => Buffer.from(await response.arrayBuffer()), + return ReadUrlResponseFactory.fromNodeJSReadable(response.body, { etag: response.headers.get('ETag') ?? undefined, - }; + }); } const message = `${url} could not be read as ${builtUrl}, ${response.status} ${response.statusText}`; diff --git a/packages/backend-common/src/reading/GoogleGcsUrlReader.ts b/packages/backend-common/src/reading/GoogleGcsUrlReader.ts index baa9477c6a..d2c3e96766 100644 --- a/packages/backend-common/src/reading/GoogleGcsUrlReader.ts +++ b/packages/backend-common/src/reading/GoogleGcsUrlReader.ts @@ -28,6 +28,8 @@ import { GoogleGcsIntegrationConfig, readGoogleGcsIntegrationConfig, } from '@backstage/integration'; +import { Readable } from 'stream'; +import { ReadUrlResponseFactory } from './ReadUrlResponseFactory'; const GOOGLE_GCS_HOST = 'storage.cloud.google.com'; @@ -85,13 +87,14 @@ export class GoogleGcsUrlReader implements UrlReader { private readonly storage: Storage, ) {} + private readStreamFromUrl(url: string): Readable { + const { bucket, key } = parseURL(url); + return this.storage.bucket(bucket).file(key).createReadStream(); + } + async read(url: string): Promise { try { - const { bucket, key } = parseURL(url); - - return await getRawBody( - this.storage.bucket(bucket).file(key).createReadStream(), - ); + return await getRawBody(this.readStreamFromUrl(url)); } catch (error) { throw new Error(`unable to read gcs file from ${url}, ${error}`); } @@ -102,8 +105,8 @@ export class GoogleGcsUrlReader implements UrlReader { _options?: ReadUrlOptions, ): Promise { // TODO etag is not implemented yet. - const buffer = await this.read(url); - return { buffer: async () => buffer }; + const stream = this.readStreamFromUrl(url); + return ReadUrlResponseFactory.fromReadable(stream); } async readTree(): Promise { diff --git a/packages/backend-common/src/reading/ReadUrlResponseFactory.test.ts b/packages/backend-common/src/reading/ReadUrlResponseFactory.test.ts new file mode 100644 index 0000000000..c8a659e792 --- /dev/null +++ b/packages/backend-common/src/reading/ReadUrlResponseFactory.test.ts @@ -0,0 +1,125 @@ +/* + * Copyright 2022 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 { ConflictError } from '@backstage/errors'; +import getRawBody from 'raw-body'; +import { Readable, Stream } from 'stream'; +import { ReadUrlResponseFactory } from './ReadUrlResponseFactory'; + +describe('ReadUrlResponseFactory', () => { + describe("fromReadable's", () => { + const expectedText = 'expected text'; + let readable: Readable; + + beforeEach(() => { + readable = Readable.from(expectedText, { + objectMode: false, + }); + }); + + it('etag is passed through', async () => { + const expectedEtag = 'xyz'; + const response = await ReadUrlResponseFactory.fromReadable(readable, { + etag: expectedEtag, + }); + expect(response.etag).toEqual(expectedEtag); + }); + + it('buffer returns expected data', async () => { + const response = await ReadUrlResponseFactory.fromReadable(readable); + const buffer = await response.buffer(); + expect(buffer.toString()).toEqual(expectedText); + }); + + it('buffer can be called multiple times', async () => { + const response = await ReadUrlResponseFactory.fromReadable(readable); + const buffer1 = await response.buffer(); + const buffer2 = await response.buffer(); + expect(buffer1.toString()).toEqual(expectedText); + expect(buffer2.toString()).toEqual(expectedText); + }); + + it('buffer cannot be called after stream is called', async () => { + const response = await ReadUrlResponseFactory.fromReadable(readable); + response.stream!(); + expect(() => response.buffer()).toThrowError(ConflictError); + }); + + it('stream returns expected data', async () => { + const response = await ReadUrlResponseFactory.fromReadable(readable); + const stream = response.stream!(); + const bufferFromStream = await getRawBody(stream); + expect(bufferFromStream.toString()).toEqual(expectedText); + }); + + it('stream cannot be called after buffer is called', async () => { + const response = await ReadUrlResponseFactory.fromReadable(readable); + response.buffer(); + expect(() => response.stream!()).toThrowError(ConflictError); + }); + }); + + describe("fromNodeJSReadable's", () => { + const expectedText = 'expected text'; + const timeouts: NodeJS.Timeout[] = []; + let readable: NodeJS.ReadableStream; + + beforeEach(() => { + // @ts-ignore + readable = new Stream({ encoding: 'utf-8' }); + readable.readable = true; + + // Write data asynchronously, as soon as possible. + timeouts[0] = setTimeout(() => { + timeouts[1] = setTimeout(readable.emit.bind(readable, 'end'), 0); + readable.emit('data', expectedText); + }, 0); + }); + + afterEach(() => { + // Clear out timeouts so we don't emit data across tests. + timeouts.forEach(clearTimeout); + }); + + it('etag is passed through', async () => { + const expectedEtag = 'xyz'; + const response = await ReadUrlResponseFactory.fromNodeJSReadable( + readable, + { + etag: expectedEtag, + }, + ); + expect(response.etag).toEqual(expectedEtag); + }); + + it('buffer returns expected data', async () => { + const response = await ReadUrlResponseFactory.fromNodeJSReadable( + readable, + ); + const buffer = await response.buffer(); + expect(buffer.toString()).toEqual(expectedText); + }); + + it('stream returns expected data', async () => { + const response = await ReadUrlResponseFactory.fromNodeJSReadable( + readable, + ); + const stream = response.stream!(); + const bufferFromStream = await getRawBody(stream); + expect(bufferFromStream.toString()).toEqual(expectedText); + }); + }); +}); diff --git a/packages/backend-common/src/reading/ReadUrlResponseFactory.ts b/packages/backend-common/src/reading/ReadUrlResponseFactory.ts new file mode 100644 index 0000000000..91d1939c24 --- /dev/null +++ b/packages/backend-common/src/reading/ReadUrlResponseFactory.ts @@ -0,0 +1,77 @@ +/* + * Copyright 2022 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 { ConflictError } from '@backstage/errors'; +import getRawBody from 'raw-body'; +import { Readable } from 'stream'; +import { + ReadUrlResponse, + ReadUrlResponseFactoryFromStreamOptions, +} from './types'; + +/** + * Utility class for UrlReader implementations to create valid ReadUrlResponse + * instances from common response primitives. + * + * @public + */ +export class ReadUrlResponseFactory { + /** + * Resolves a ReadUrlResponse from a Readable stream. + */ + static async fromReadable( + stream: Readable, + options?: ReadUrlResponseFactoryFromStreamOptions, + ): Promise { + // Reference to eventual buffer enables callers to call buffer() multiple + // times without consequence. + let buffer: Promise; + + // Prevent "stream is not readable" errors from bubbling up. + const conflictError = new ConflictError( + 'Cannot use buffer() and stream() from the same ReadUrlResponse', + ); + let hasCalledStream = false; + let hasCalledBuffer = false; + + return { + buffer: () => { + hasCalledBuffer = true; + if (hasCalledStream) throw conflictError; + if (buffer) return buffer; + buffer = getRawBody(stream); + return buffer; + }, + stream: () => { + hasCalledStream = true; + if (hasCalledBuffer) throw conflictError; + return stream; + }, + etag: options?.etag, + }; + } + + /** + * Resolves a ReadUrlResponse from an old-style NodeJS.ReadableStream. + */ + static async fromNodeJSReadable( + oldStyleStream: NodeJS.ReadableStream, + options?: ReadUrlResponseFactoryFromStreamOptions, + ): Promise { + const readable = new Readable().wrap(oldStyleStream); + return ReadUrlResponseFactory.fromReadable(readable, options); + } +} diff --git a/packages/backend-common/src/reading/index.ts b/packages/backend-common/src/reading/index.ts index f1609179b5..9ee86ef3f2 100644 --- a/packages/backend-common/src/reading/index.ts +++ b/packages/backend-common/src/reading/index.ts @@ -23,6 +23,7 @@ export { GithubUrlReader } from './GithubUrlReader'; export { GitlabUrlReader } from './GitlabUrlReader'; export { AwsS3UrlReader } from './AwsS3UrlReader'; export { FetchUrlReader } from './FetchUrlReader'; +export { ReadUrlResponseFactory } from './ReadUrlResponseFactory'; export type { FromReadableArrayOptions, ReaderFactory, @@ -34,6 +35,7 @@ export type { ReadTreeResponseFactoryOptions, ReadUrlOptions, ReadUrlResponse, + ReadUrlResponseFactoryFromStreamOptions, SearchOptions, SearchResponse, SearchResponseFile, diff --git a/packages/backend-common/src/reading/types.ts b/packages/backend-common/src/reading/types.ts index 9505778705..6eae4e6cac 100644 --- a/packages/backend-common/src/reading/types.ts +++ b/packages/backend-common/src/reading/types.ts @@ -124,6 +124,15 @@ export type ReadUrlResponse = { */ buffer(): Promise; + /** + * Returns the data that was read from the remote URL as a Readable stream. + * + * @remarks + * + * This method will be required in a future release. + */ + stream?(): Readable; + /** * Etag returned by content provider. * @@ -134,6 +143,15 @@ export type ReadUrlResponse = { etag?: string; }; +/** + * An options object for {@link ReadUrlResponseFactory} factory methods. + * + * @public + */ +export type ReadUrlResponseFactoryFromStreamOptions = { + etag?: string; +}; + /** * An options object for {@link UrlReader.readTree} operations. * diff --git a/yarn.lock b/yarn.lock index 24dfd6a6c0..6de5d62b97 100644 --- a/yarn.lock +++ b/yarn.lock @@ -5417,6 +5417,13 @@ dependencies: "@babel/types" "^7.3.0" +"@types/base64-stream@^1.0.2": + version "1.0.2" + resolved "https://registry.npmjs.org/@types/base64-stream/-/base64-stream-1.0.2.tgz#5a6c99eb4701a3d2aafd3082d8b44be7b1be1325" + integrity sha512-pPTps4oDwasxGcxpG6Ub8rAkVSPdekoolze0oza1mMFdlgXXKXE2zqCePBfemvWlGI92k9jE7kH3nz0+tJEXQw== + dependencies: + "@types/node" "*" + "@types/body-parser@*", "@types/body-parser@1.19.2", "@types/body-parser@^1.19.0": version "1.19.2" resolved "https://registry.npmjs.org/@types/body-parser/-/body-parser-1.19.2.tgz#aea2059e28b7658639081347ac4fab3de166e6f0" @@ -7979,6 +7986,11 @@ base64-js@^1.0.2, base64-js@^1.3.0, base64-js@^1.3.1, base64-js@^1.5.1: resolved "https://registry.npmjs.org/base64-js/-/base64-js-1.5.1.tgz#1b1b440160a5bf7ad40b650f095963481903930a" integrity sha512-AKpaYlHn8t4SVbOHCy+b5+KKgvR4vrsD8vbvrbiQJps7fKDTkjkDry6ji0rUJjC0kzbNePLwzxq8iypo41qeWA== +base64-stream@^1.0.0: + version "1.0.0" + resolved "https://registry.npmjs.org/base64-stream/-/base64-stream-1.0.0.tgz#157ae00bc7888695e884e1fcc51c551fdfa8a1fa" + integrity sha512-BQQZftaO48FcE1Kof9CmXMFaAdqkcNorgc8CxesZv9nMbbTF1EFyQe89UOuh//QMmdtfUDXyO8rgUalemL5ODA== + base64id@2.0.0: version "2.0.0" resolved "https://registry.npmjs.org/base64id/-/base64id-2.0.0.tgz#2770ac6bc47d312af97a8bf9a634342e0cd25cb6"