Add stream() to UrlReader.readUrl response
Signed-off-by: Eric Peterson <ericpeterson@spotify.com>
This commit is contained in:
@@ -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.
|
||||
@@ -547,6 +547,24 @@ export type ReadUrlOptions = {
|
||||
// @public
|
||||
export type ReadUrlResponse = {
|
||||
buffer(): Promise<Buffer>;
|
||||
stream?(): Readable;
|
||||
etag?: string;
|
||||
};
|
||||
|
||||
// @public
|
||||
export class ReadUrlResponseFactory {
|
||||
static fromNodeJSReadable(
|
||||
oldStyleStream: NodeJS.ReadableStream,
|
||||
options?: ReadUrlResponseFactoryFromStreamOptions,
|
||||
): Promise<ReadUrlResponse>;
|
||||
static fromReadable(
|
||||
stream: Readable,
|
||||
options?: ReadUrlResponseFactoryFromStreamOptions,
|
||||
): Promise<ReadUrlResponse>;
|
||||
}
|
||||
|
||||
// @public
|
||||
export type ReadUrlResponseFactoryFromStreamOptions = {
|
||||
etag?: string;
|
||||
};
|
||||
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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}`;
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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}`;
|
||||
|
||||
@@ -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}`;
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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}`;
|
||||
|
||||
@@ -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}`;
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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}`;
|
||||
|
||||
@@ -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}`;
|
||||
|
||||
@@ -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<Buffer> {
|
||||
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<ReadUrlResponse> {
|
||||
// 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<ReadTreeResponse> {
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
});
|
||||
@@ -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<ReadUrlResponse> {
|
||||
// Reference to eventual buffer enables callers to call buffer() multiple
|
||||
// times without consequence.
|
||||
let buffer: Promise<Buffer>;
|
||||
|
||||
// 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<ReadUrlResponse> {
|
||||
const readable = new Readable().wrap(oldStyleStream);
|
||||
return ReadUrlResponseFactory.fromReadable(readable, options);
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
|
||||
@@ -124,6 +124,15 @@ export type ReadUrlResponse = {
|
||||
*/
|
||||
buffer(): Promise<Buffer>;
|
||||
|
||||
/**
|
||||
* 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.
|
||||
*
|
||||
|
||||
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user