diff --git a/packages/sdk/electron/__tests__/platform/ElectronRequests.streaming.test.ts b/packages/sdk/electron/__tests__/platform/ElectronRequests.streaming.test.ts new file mode 100644 index 0000000000..e19bad3cc7 --- /dev/null +++ b/packages/sdk/electron/__tests__/platform/ElectronRequests.streaming.test.ts @@ -0,0 +1,149 @@ +import { + AsyncQueue, + TestHttpHandlers, + TestHttpServer, + TestHttpServers, +} from 'launchdarkly-js-test-helpers'; + +import ElectronRequests from '../../src/platform/ElectronRequests'; + +describe('given a running HTTP server', () => { + let server: TestHttpServer; + + beforeEach(async () => { + server = await TestHttpServers.start(); + }); + + afterEach(async () => { + await server.closeAndWait(); + }); + + it('forwards a streaming request and exposes the status, headers, and body chunks', async () => { + const chunks = new AsyncQueue(); + chunks.add('first'); + server.byDefault( + TestHttpHandlers.chunkedStream(200, { 'content-type': 'text/event-stream' }, chunks), + ); + + const requests = new ElectronRequests(); + const res = await requests.fetch(`${server.url}/stream`, { + method: 'REPORT', + headers: { authorization: 'sdk-key' }, + body: '{"kind":"user"}', + streaming: true, + }); + + expect(res.status).toEqual(200); + const collected: Record = {}; + res.headers.forEach?.((value, key) => { + collected[key] = value; + }); + expect(collected['content-type']).toEqual('text/event-stream'); + + const reader = res.body?.getReader(); + const first = await reader?.read(); + expect(first?.done).toBe(false); + expect(Buffer.from(first?.value ?? []).toString()).toEqual('first'); + + const received = await server.nextRequest(); + expect(received.method.toUpperCase()).toEqual('REPORT'); + expect(received.headers.authorization).toEqual('sdk-key'); + expect(received.body).toEqual('{"kind":"user"}'); + }); + + it('does not request compressed content for a streaming request', async () => { + const chunks = new AsyncQueue(); + server.byDefault(TestHttpHandlers.chunkedStream(200, {}, chunks)); + + const requests = new ElectronRequests(); + await requests.fetch(server.url, { method: 'GET', streaming: true }); + + const received = await server.nextRequest(); + expect(received.headers['accept-encoding']).toBeUndefined(); + }); + + it('does not follow redirects for a streaming request', async () => { + server.byDefault(TestHttpHandlers.respond(301, { location: `${server.url}/other` })); + + const requests = new ElectronRequests(); + const res = await requests.fetch(server.url, { method: 'GET', streaming: true }); + + expect(res.status).toEqual(301); + expect(server.requestCount()).toEqual(1); + }); + + it('stops the stream when the signal aborts', async () => { + const chunks = new AsyncQueue(); + chunks.add('first'); + server.byDefault(TestHttpHandlers.chunkedStream(200, {}, chunks)); + + const controller = new AbortController(); + const requests = new ElectronRequests(); + const res = await requests.fetch(server.url, { + method: 'GET', + streaming: true, + signal: controller.signal, + }); + const reader = res.body?.getReader(); + await reader?.read(); + const pending = reader?.read(); + controller.abort(); + await expect(pending).rejects.toThrow(); + }); + + it('rejects a streaming request when the signal is already aborted', async () => { + server.byDefault(TestHttpHandlers.respond(200)); + const controller = new AbortController(); + controller.abort(); + + const requests = new ElectronRequests(); + await expect( + requests.fetch(server.url, { method: 'GET', streaming: true, signal: controller.signal }), + ).rejects.toThrow(); + }); + + it('streams SSE events through createEventSource', async () => { + const chunks = new AsyncQueue(); + chunks.add('data: hello\n\n'); + server.byDefault( + TestHttpHandlers.chunkedStream(200, { 'content-type': 'text/event-stream' }, chunks), + ); + + const requests = new ElectronRequests(); + const es = requests.createEventSource(`${server.url}/stream`, { + headers: {}, + initialRetryDelayMillis: 100, + readTimeoutMillis: 5000, + retryResetIntervalMillis: 30_000, + errorFilter: () => false, + }); + try { + const messages = new AsyncQueue<{ data?: string }>(); + es.addEventListener('message', (event) => messages.add(event ?? {})); + const message = await messages.take(); + expect(message.data).toEqual('hello'); + } finally { + es.close(); + } + }); +}); + +describe('given a running HTTPS server with a self-signed certificate', () => { + let server: TestHttpServer; + + beforeEach(async () => { + server = await TestHttpServers.startSecure(); + server.byDefault(TestHttpHandlers.respond(200)); + }); + + afterEach(async () => { + await server.closeAndWait(); + }); + + it('rejects a streaming request when the certificate is not trusted', async () => { + // This SDK exposes no TLS options. Verification follows the platform default, so a + // self-signed certificate the machine does not trust must fail the request. + const requests = new ElectronRequests(); + await expect(requests.fetch(server.url, { method: 'GET', streaming: true })).rejects.toThrow(); + }); +}); diff --git a/packages/sdk/electron/__tests__/platform/HeaderWrapper.test.ts b/packages/sdk/electron/__tests__/platform/HeaderWrapper.test.ts index d2c89ae6b9..bc1f40f1b8 100644 --- a/packages/sdk/electron/__tests__/platform/HeaderWrapper.test.ts +++ b/packages/sdk/electron/__tests__/platform/HeaderWrapper.test.ts @@ -50,4 +50,16 @@ describe('given header values', () => { } expect(values).toEqual(['anything', 'some-value', 'a, b']); }); + + it('iterates each header with the value before the key', () => { + const collected: [string, string][] = []; + wrapper.forEach((value, key) => { + collected.push([value, key]); + }); + expect(collected).toEqual([ + ['anything', 'accept'], + ['some-value', 'some-header'], + ['a, b', 'some-array'], + ]); + }); }); diff --git a/packages/sdk/electron/package.json b/packages/sdk/electron/package.json index ae38ae5d5e..8556b45738 100644 --- a/packages/sdk/electron/package.json +++ b/packages/sdk/electron/package.json @@ -53,13 +53,14 @@ "electron": ">=34.5.6" }, "dependencies": { - "@launchdarkly/js-client-sdk-common": "workspace:^", - "launchdarkly-eventsource": "2.2.0" + "@launchdarkly/eventsource": "workspace:^", + "@launchdarkly/js-client-sdk-common": "workspace:^" }, "devDependencies": { "@types/jest": "^29.4.0", "electron": "^42.4.0", "jest": "^29.5.0", + "launchdarkly-js-test-helpers": "^2.2.0", "oxfmt": "0.63.0", "oxlint": "1.78.0", "ts-jest": "^29.0.5", diff --git a/packages/sdk/electron/src/platform/ElectronRequests.ts b/packages/sdk/electron/src/platform/ElectronRequests.ts index 18a3fd7712..2e05928c58 100644 --- a/packages/sdk/electron/src/platform/ElectronRequests.ts +++ b/packages/sdk/electron/src/platform/ElectronRequests.ts @@ -1,12 +1,10 @@ import * as http from 'http'; import * as https from 'https'; -// No types for the event source. -// @ts-ignore -import { EventSource as LDEventSource } from 'launchdarkly-eventsource'; import { promisify } from 'util'; import * as zlib from 'zlib'; -import { EventSourceCapabilities, platform } from '@launchdarkly/js-client-sdk-common'; +import { createEventSource } from '@launchdarkly/eventsource'; +import { EventSourceCapabilities, internal, platform } from '@launchdarkly/js-client-sdk-common'; import ElectronResponse from './ElectronResponse'; @@ -20,6 +18,9 @@ export default class ElectronRequests implements platform.Requests { } async fetch(url: string, options: platform.Options = {}): Promise { + if (options.streaming) { + return this._streamingFetch(url, options); + } const isSecure = url.startsWith('https://'); const impl = isSecure ? https : http; @@ -27,7 +28,7 @@ export default class ElectronRequests implements platform.Requests { let bodyData: string | Buffer | undefined = options.body; // For get requests we are going to automatically support compressed responses. - // Note this does not affect SSE as the event source is not using this fetch implementation. + // Note this does not affect SSE as streaming requests take the branch above. if (options.method?.toLowerCase() === 'get') { headers['accept-encoding'] = 'gzip'; } @@ -67,16 +68,56 @@ export default class ElectronRequests implements platform.Requests { }); } + /** + * The transport for a streaming request. This SDK exposes no agent, proxy, or TLS options, so + * the running machine's own network configuration applies and TLS verification follows the + * platform default. It does not request compressed content, and it never follows a redirect. + * A redirect status resolves like any other non-200 response, and the caller decides whether + * to retry the original URL. It applies no read or socket timeout. The caller owns the read + * timeout and cancels through the abort signal. + */ + private _streamingFetch(url: string, options: platform.Options): Promise { + const isSecure = url.startsWith('https://'); + const impl = isSecure ? https : http; + const requestOptions: https.RequestOptions = { + method: options.method, + headers: options.headers, + }; + return new Promise((resolve, reject) => { + const req = impl.request(url, requestOptions, (res) => + resolve(internal.createStreamingResponse(res)), + ); + // An SSE consumer wants each chunk as soon as it arrives; do not batch small writes. + req.setNoDelay(true); + const { signal } = options; + if (signal) { + const abort = () => req.destroy(new Error('The stream request was aborted')); + if (signal.aborted) { + abort(); + } else { + signal.addEventListener('abort', abort, { once: true }); + } + } + // This listener stays attached after resolve. A later socket error then becomes a harmless + // no-op reject instead of an unhandled 'error' event that would crash the process + req.on('error', reject); + if (options.body !== undefined) { + req.write(options.body); + } + req.end(); + }); + } + createEventSource( url: string, eventSourceInitDict: platform.EventSourceInitDict, ): platform.EventSource { - const expandedOptions = { + return createEventSource(url, { ...eventSourceInitDict, maxBackoffMillis: 30 * 1000, jitterRatio: 0.5, - }; - return new LDEventSource(url, expandedOptions); + fetch: (fetchUrl, init) => this.fetch(fetchUrl, { ...init, streaming: true }), + }); } getEventSourceCapabilities(): EventSourceCapabilities { diff --git a/packages/sdk/electron/src/platform/HeaderWrapper.ts b/packages/sdk/electron/src/platform/HeaderWrapper.ts index 9c447c2a8f..f6922bf30d 100644 --- a/packages/sdk/electron/src/platform/HeaderWrapper.ts +++ b/packages/sdk/electron/src/platform/HeaderWrapper.ts @@ -52,6 +52,17 @@ export default class HeaderWrapper implements platform.Headers { } } + /** + * Executes the callback once for each header, with the value first. The order matches the + * fetch `Headers.forEach` signature. Multi-value headers are joined with a comma, and + * headers without a value are skipped, like `entries`. + */ + forEach(callback: (value: string, key: string) => void): void { + for (const [key, value] of this.entries()) { + callback(value, key); + } + } + has(name: string): boolean { return Object.prototype.hasOwnProperty.call(this._headers, name); } diff --git a/packages/sdk/node-client/__tests__/platform/HeaderWrapper.test.ts b/packages/sdk/node-client/__tests__/platform/HeaderWrapper.test.ts index 72233ae5b8..e55624f382 100644 --- a/packages/sdk/node-client/__tests__/platform/HeaderWrapper.test.ts +++ b/packages/sdk/node-client/__tests__/platform/HeaderWrapper.test.ts @@ -49,3 +49,18 @@ it('reports presence with has()', () => { expect(wrapper.has('content-type')).toBe(true); expect(wrapper.has('missing')).toBe(false); }); + +it('iterates each header with the value before the key', () => { + const wrapper = new HeaderWrapper({ + accept: 'anything', + 'some-array': ['a', 'b'], + }); + const collected: [string, string][] = []; + wrapper.forEach((value, key) => { + collected.push([value, key]); + }); + expect(collected).toEqual([ + ['anything', 'accept'], + ['a, b', 'some-array'], + ]); +}); diff --git a/packages/sdk/node-client/__tests__/platform/NodeRequests.streaming.test.ts b/packages/sdk/node-client/__tests__/platform/NodeRequests.streaming.test.ts new file mode 100644 index 0000000000..135d5ef3d2 --- /dev/null +++ b/packages/sdk/node-client/__tests__/platform/NodeRequests.streaming.test.ts @@ -0,0 +1,206 @@ +import http from 'http'; +import https from 'https'; +import { + AsyncQueue, + TestHttpHandlers, + TestHttpServer, + TestHttpServers, +} from 'launchdarkly-js-test-helpers'; + +import NodeRequests from '../../src/platform/NodeRequests'; + +describe('given a running HTTP server', () => { + let server: TestHttpServer; + + beforeEach(async () => { + server = await TestHttpServers.start(); + }); + + afterEach(async () => { + jest.restoreAllMocks(); + await server.closeAndWait(); + }); + + it('forwards a streaming request and exposes the status, headers, and body chunks', async () => { + const chunks = new AsyncQueue(); + chunks.add('first'); + server.byDefault( + TestHttpHandlers.chunkedStream(200, { 'content-type': 'text/event-stream' }, chunks), + ); + + const requests = new NodeRequests(); + const res = await requests.fetch(`${server.url}/stream`, { + method: 'REPORT', + headers: { authorization: 'sdk-key' }, + body: '{"kind":"user"}', + streaming: true, + }); + + expect(res.status).toEqual(200); + const collected: Record = {}; + res.headers.forEach?.((value, key) => { + collected[key] = value; + }); + expect(collected['content-type']).toEqual('text/event-stream'); + + const reader = res.body?.getReader(); + const first = await reader?.read(); + expect(first?.done).toBe(false); + expect(Buffer.from(first?.value ?? []).toString()).toEqual('first'); + + const received = await server.nextRequest(); + expect(received.method.toUpperCase()).toEqual('REPORT'); + expect(received.headers.authorization).toEqual('sdk-key'); + expect(received.body).toEqual('{"kind":"user"}'); + }); + + it('does not request compressed content for a streaming request', async () => { + const chunks = new AsyncQueue(); + server.byDefault(TestHttpHandlers.chunkedStream(200, {}, chunks)); + + const requests = new NodeRequests(); + await requests.fetch(server.url, { method: 'GET', streaming: true }); + + const received = await server.nextRequest(); + expect(received.headers['accept-encoding']).toBeUndefined(); + }); + + it('does not follow redirects for a streaming request', async () => { + server.byDefault(TestHttpHandlers.respond(301, { location: `${server.url}/other` })); + + const requests = new NodeRequests(); + const res = await requests.fetch(server.url, { method: 'GET', streaming: true }); + + expect(res.status).toEqual(301); + expect(server.requestCount()).toEqual(1); + }); + + it('stops the stream when the signal aborts', async () => { + const chunks = new AsyncQueue(); + chunks.add('first'); + server.byDefault(TestHttpHandlers.chunkedStream(200, {}, chunks)); + + const controller = new AbortController(); + const requests = new NodeRequests(); + const res = await requests.fetch(server.url, { + method: 'GET', + streaming: true, + signal: controller.signal, + }); + const reader = res.body?.getReader(); + await reader?.read(); + const pending = reader?.read(); + controller.abort(); + await expect(pending).rejects.toThrow(); + }); + + it('rejects a streaming request when the signal is already aborted', async () => { + server.byDefault(TestHttpHandlers.respond(200)); + const controller = new AbortController(); + controller.abort(); + + const requests = new NodeRequests(); + await expect( + requests.fetch(server.url, { method: 'GET', streaming: true, signal: controller.signal }), + ).rejects.toThrow(); + }); + + it('does not pass the CA to the per-request options of an http streaming request', async () => { + const requestSpy = jest.spyOn(http, 'request'); + + const requests = new NodeRequests({ ca: 'unused-ca' }); + // The https agent built from the TLS options is rejected for an http URL, so the request + // fails after http.request has already received its options. + await requests.fetch(server.url, { method: 'GET', streaming: true }).catch(() => undefined); + + expect(requestSpy).toHaveBeenCalledTimes(1); + expect(requestSpy.mock.calls[0][1]).not.toHaveProperty('ca'); + }); + + it('streams SSE events through createEventSource', async () => { + const chunks = new AsyncQueue(); + chunks.add('data: hello\n\n'); + server.byDefault( + TestHttpHandlers.chunkedStream(200, { 'content-type': 'text/event-stream' }, chunks), + ); + + const requests = new NodeRequests(); + const es = requests.createEventSource(`${server.url}/stream`, { + headers: {}, + initialRetryDelayMillis: 100, + readTimeoutMillis: 5000, + retryResetIntervalMillis: 30_000, + errorFilter: () => false, + }); + try { + const messages = new AsyncQueue<{ data?: string }>(); + es.addEventListener('message', (event) => messages.add(event ?? {})); + const message = await messages.take(); + expect(message.data).toEqual('hello'); + } finally { + es.close(); + } + }); +}); + +describe('given a running HTTPS server with a self-signed certificate', () => { + let server: TestHttpServer; + + beforeEach(async () => { + server = await TestHttpServers.startSecure(); + server.byDefault(TestHttpHandlers.respond(200)); + }); + + afterEach(async () => { + jest.restoreAllMocks(); + await server.closeAndWait(); + }); + + it('connects a streaming request when the CA is supplied in the TLS options', async () => { + const requests = new NodeRequests({ ca: server.certificate }); + const res = await requests.fetch(server.url, { method: 'GET', streaming: true }); + expect(res.status).toEqual(200); + }); + + it('passes the CA to the per-request options of an https streaming request', async () => { + const requestSpy = jest.spyOn(https, 'request'); + + const requests = new NodeRequests({ ca: server.certificate }); + await requests.fetch(server.url, { method: 'GET', streaming: true }); + + expect(requestSpy).toHaveBeenCalledTimes(1); + expect(requestSpy.mock.calls[0][1]).toEqual( + expect.objectContaining({ ca: server.certificate }), + ); + }); + + it('rejects a streaming request when no CA is supplied', async () => { + const requests = new NodeRequests(); + await expect(requests.fetch(server.url, { method: 'GET', streaming: true })).rejects.toThrow(); + }); + + it('streams over TLS when the CA is supplied in the TLS options', async () => { + const chunks = new AsyncQueue(); + chunks.add('data: secure\n\n'); + server.byDefault( + TestHttpHandlers.chunkedStream(200, { 'content-type': 'text/event-stream' }, chunks), + ); + + const requests = new NodeRequests({ ca: server.certificate }); + const es = requests.createEventSource(`${server.url}/stream`, { + headers: {}, + initialRetryDelayMillis: 100, + readTimeoutMillis: 5000, + retryResetIntervalMillis: 30_000, + errorFilter: () => false, + }); + try { + const messages = new AsyncQueue<{ data?: string }>(); + es.addEventListener('message', (event) => messages.add(event ?? {})); + const message = await messages.take(); + expect(message.data).toEqual('secure'); + } finally { + es.close(); + } + }); +}); diff --git a/packages/sdk/node-client/package.json b/packages/sdk/node-client/package.json index cfa0f6aa7f..589239807b 100644 --- a/packages/sdk/node-client/package.json +++ b/packages/sdk/node-client/package.json @@ -44,8 +44,8 @@ "contract-tests": "./contract-tests/run-contract-tests.sh" }, "dependencies": { - "@launchdarkly/js-client-sdk-common": "1.32.5", - "launchdarkly-eventsource": "2.2.0" + "@launchdarkly/eventsource": "0.1.0", + "@launchdarkly/js-client-sdk-common": "1.32.5" }, "devDependencies": { "@types/jest": "^29.4.0", diff --git a/packages/sdk/node-client/src/platform/HeaderWrapper.ts b/packages/sdk/node-client/src/platform/HeaderWrapper.ts index d3812f1937..9b30516989 100644 --- a/packages/sdk/node-client/src/platform/HeaderWrapper.ts +++ b/packages/sdk/node-client/src/platform/HeaderWrapper.ts @@ -50,6 +50,17 @@ export default class HeaderWrapper implements platform.Headers { } } + /** + * Executes the callback once for each header, with the value first. The order matches the + * fetch `Headers.forEach` signature. Multi-value headers are joined with a comma, and + * headers without a value are skipped, like `entries`. + */ + forEach(callback: (value: string, key: string) => void): void { + for (const [key, value] of this.entries()) { + callback(value, key); + } + } + has(name: string): boolean { return Object.prototype.hasOwnProperty.call(this._headers, name); } diff --git a/packages/sdk/node-client/src/platform/NodeRequests.ts b/packages/sdk/node-client/src/platform/NodeRequests.ts index 1ed819e94d..d7cefd269f 100644 --- a/packages/sdk/node-client/src/platform/NodeRequests.ts +++ b/packages/sdk/node-client/src/platform/NodeRequests.ts @@ -1,18 +1,47 @@ import * as http from 'http'; import * as https from 'https'; -// No types for the event source. -// @ts-ignore -import { EventSource as LDEventSource } from 'launchdarkly-eventsource'; import { promisify } from 'util'; import * as zlib from 'zlib'; -import { EventSourceCapabilities, platform } from '@launchdarkly/js-client-sdk-common'; +import { createEventSource } from '@launchdarkly/eventsource'; +import { EventSourceCapabilities, internal, platform } from '@launchdarkly/js-client-sdk-common'; import type { LDTLSOptions } from '../NodeOptions'; import NodeResponse from './NodeResponse'; const gzip = promisify(zlib.gzip); +/** + * The TLS options that are copied onto each streaming `https` request. The names match the + * options of `https.request()`. + */ +const TLS_OPTION_NAMES = [ + 'pfx', + 'key', + 'passphrase', + 'cert', + 'ca', + 'ciphers', + 'rejectUnauthorized', + 'secureProtocol', + 'servername', + 'checkServerIdentity', +] as const; + +function tlsRequestOptions(tlsOptions?: LDTLSOptions): Record { + const merged: Record = {}; + if (!tlsOptions) { + return merged; + } + const bag = tlsOptions as unknown as Record; + TLS_OPTION_NAMES.forEach((name) => { + if (bag[name] !== undefined) { + merged[name] = bag[name]; + } + }); + return merged; +} + function processTlsOptions(tlsOptions: LDTLSOptions): https.AgentOptions { const options: https.AgentOptions & { [index: string]: any } = { ca: tlsOptions.ca, @@ -43,17 +72,20 @@ function processTlsOptions(tlsOptions: LDTLSOptions): https.AgentOptions { export default class NodeRequests implements platform.Requests { private _agent: https.Agent | undefined; - private _tlsOptions: LDTLSOptions | undefined; + private _tlsParams: Record; private _enableBodyCompression: boolean = false; constructor(tlsOptions?: LDTLSOptions, enableEventCompression?: boolean) { this._agent = tlsOptions ? new https.Agent(processTlsOptions(tlsOptions)) : undefined; - this._tlsOptions = tlsOptions; + this._tlsParams = tlsRequestOptions(tlsOptions); this._enableBodyCompression = !!enableEventCompression; } async fetch(url: string, options: platform.Options = {}): Promise { + if (options.streaming) { + return this._streamingFetch(url, options); + } const isSecure = url.startsWith('https://'); const impl = isSecure ? https : http; @@ -100,18 +132,58 @@ export default class NodeRequests implements platform.Requests { }); } + /** + * The transport for a streaming request. It does not request compressed content, and it never + * follows a redirect. A redirect status resolves like any other non-200 response, and the + * caller decides whether to retry the original URL. It applies no read or socket timeout. The + * caller owns the read timeout and cancels through the abort signal. + */ + private _streamingFetch(url: string, options: platform.Options): Promise { + const isSecure = url.startsWith('https://'); + const impl = isSecure ? https : http; + const requestOptions: https.RequestOptions & Record = { + method: options.method, + headers: options.headers, + agent: this._agent, + }; + if (isSecure) { + Object.assign(requestOptions, this._tlsParams); + } + return new Promise((resolve, reject) => { + const req = impl.request(url, requestOptions, (res) => + resolve(internal.createStreamingResponse(res)), + ); + // An SSE consumer wants each chunk as soon as it arrives; do not batch small writes. + req.setNoDelay(true); + const { signal } = options; + if (signal) { + const abort = () => req.destroy(new Error('The stream request was aborted')); + if (signal.aborted) { + abort(); + } else { + signal.addEventListener('abort', abort, { once: true }); + } + } + // This listener stays attached after resolve. A later socket error then becomes a harmless + // no-op reject instead of an unhandled 'error' event that would crash the process + req.on('error', reject); + if (options.body !== undefined) { + req.write(options.body); + } + req.end(); + }); + } + createEventSource( url: string, eventSourceInitDict: platform.EventSourceInitDict, ): platform.EventSource { - const expandedOptions = { + return createEventSource(url, { ...eventSourceInitDict, - agent: this._agent, - tlsParams: this._tlsOptions, maxBackoffMillis: 30 * 1000, jitterRatio: 0.5, - }; - return new LDEventSource(url, expandedOptions); + fetch: (fetchUrl, init) => this.fetch(fetchUrl, { ...init, streaming: true }), + }); } getEventSourceCapabilities(): EventSourceCapabilities { diff --git a/packages/sdk/server-node/__tests__/platform/HeaderWrapper.test.ts b/packages/sdk/server-node/__tests__/platform/HeaderWrapper.test.ts index bf1cfd0b50..680cefc3c5 100644 --- a/packages/sdk/server-node/__tests__/platform/HeaderWrapper.test.ts +++ b/packages/sdk/server-node/__tests__/platform/HeaderWrapper.test.ts @@ -50,4 +50,16 @@ describe('given header values', () => { } expect(values).toEqual(['anything', 'some-value', 'a, b']); }); + + it('iterates each header with the value before the key', () => { + const collected: [string, string][] = []; + wrapper.forEach((value, key) => { + collected.push([value, key]); + }); + expect(collected).toEqual([ + ['anything', 'accept'], + ['some-value', 'some-header'], + ['a, b', 'some-array'], + ]); + }); }); diff --git a/packages/sdk/server-node/__tests__/platform/NodeRequests.streaming.test.ts b/packages/sdk/server-node/__tests__/platform/NodeRequests.streaming.test.ts new file mode 100644 index 0000000000..fa2633883c --- /dev/null +++ b/packages/sdk/server-node/__tests__/platform/NodeRequests.streaming.test.ts @@ -0,0 +1,253 @@ +import * as http from 'http'; +import * as https from 'https'; +import { + AsyncQueue, + TestHttpHandlers, + TestHttpServer, + TestHttpServers, +} from 'launchdarkly-js-test-helpers'; + +import NodeRequests from '../../src/platform/NodeRequests'; + +describe('given a running HTTP server', () => { + let server: TestHttpServer; + + beforeEach(async () => { + server = await TestHttpServers.start(); + }); + + afterEach(async () => { + jest.restoreAllMocks(); + await server.closeAndWait(); + }); + + it('forwards a streaming request and exposes the status, headers, and body chunks', async () => { + const chunks = new AsyncQueue(); + chunks.add('first'); + server.byDefault( + TestHttpHandlers.chunkedStream(200, { 'content-type': 'text/event-stream' }, chunks), + ); + + const requests = new NodeRequests(); + const res = await requests.fetch(`${server.url}/stream`, { + method: 'REPORT', + headers: { authorization: 'sdk-key' }, + body: '{"kind":"user"}', + streaming: true, + }); + + expect(res.status).toEqual(200); + const collected: Record = {}; + res.headers.forEach?.((value, key) => { + collected[key] = value; + }); + expect(collected['content-type']).toEqual('text/event-stream'); + + const reader = res.body?.getReader(); + const first = await reader?.read(); + expect(first?.done).toBe(false); + expect(Buffer.from(first?.value ?? []).toString()).toEqual('first'); + + const received = await server.nextRequest(); + expect(received.method.toUpperCase()).toEqual('REPORT'); + expect(received.headers.authorization).toEqual('sdk-key'); + expect(received.body).toEqual('{"kind":"user"}'); + }); + + it('does not request compressed content for a streaming request', async () => { + const chunks = new AsyncQueue(); + server.byDefault(TestHttpHandlers.chunkedStream(200, {}, chunks)); + + const requests = new NodeRequests(); + await requests.fetch(server.url, { method: 'GET', streaming: true }); + + const received = await server.nextRequest(); + expect(received.headers['accept-encoding']).toBeUndefined(); + }); + + it('does not follow redirects for a streaming request', async () => { + server.byDefault(TestHttpHandlers.respond(301, { location: `${server.url}/other` })); + + const requests = new NodeRequests(); + const res = await requests.fetch(server.url, { method: 'GET', streaming: true }); + + expect(res.status).toEqual(301); + expect(server.requestCount()).toEqual(1); + }); + + it('stops the stream when the signal aborts', async () => { + const chunks = new AsyncQueue(); + chunks.add('first'); + server.byDefault(TestHttpHandlers.chunkedStream(200, {}, chunks)); + + const controller = new AbortController(); + const requests = new NodeRequests(); + const res = await requests.fetch(server.url, { + method: 'GET', + streaming: true, + signal: controller.signal, + }); + const reader = res.body?.getReader(); + await reader?.read(); + const pending = reader?.read(); + controller.abort(); + await expect(pending).rejects.toThrow(); + }); + + it('rejects a streaming request when the signal is already aborted', async () => { + server.byDefault(TestHttpHandlers.respond(200)); + const controller = new AbortController(); + controller.abort(); + + const requests = new NodeRequests(); + await expect( + requests.fetch(server.url, { method: 'GET', streaming: true, signal: controller.signal }), + ).rejects.toThrow(); + }); + + it('uses the supplied agent for a streaming request', async () => { + server.byDefault(TestHttpHandlers.respond(200)); + const agent = new http.Agent({ keepAlive: false }); + // addRequest exists at runtime but is not part of the public http.Agent type. + // @ts-ignore + const addRequestSpy = jest.spyOn(agent, 'addRequest'); + + const requests = new NodeRequests(undefined, undefined, agent); + await requests.fetch(server.url, { method: 'GET', streaming: true }); + + expect(addRequestSpy).toHaveBeenCalled(); + }); + + it('does not pass the CA to the per-request options of an http streaming request', async () => { + const requestSpy = jest.spyOn(http, 'request'); + + const requests = new NodeRequests({ ca: 'unused-ca' }); + // The https agent built from the TLS options is rejected for an http URL, so the request + // fails after http.request has already received its options. + await requests.fetch(server.url, { method: 'GET', streaming: true }).catch(() => undefined); + + expect(requestSpy).toHaveBeenCalledTimes(1); + expect(requestSpy.mock.calls[0][1]).not.toHaveProperty('ca'); + }); + + it('streams SSE events through createEventSource', async () => { + const chunks = new AsyncQueue(); + chunks.add('data: hello\n\n'); + server.byDefault( + TestHttpHandlers.chunkedStream(200, { 'content-type': 'text/event-stream' }, chunks), + ); + + const requests = new NodeRequests(); + const es = requests.createEventSource(`${server.url}/stream`, { + headers: {}, + initialRetryDelayMillis: 100, + readTimeoutMillis: 5000, + retryResetIntervalMillis: 30_000, + errorFilter: () => false, + }); + try { + const messages = new AsyncQueue<{ data?: string }>(); + es.addEventListener('message', (event) => messages.add(event ?? {})); + const message = await messages.take(); + expect(message.data).toEqual('hello'); + } finally { + es.close(); + } + }); +}); + +describe('given a running HTTPS server with a self-signed certificate', () => { + let server: TestHttpServer; + + beforeEach(async () => { + server = await TestHttpServers.startSecure(); + server.byDefault(TestHttpHandlers.respond(200)); + }); + + afterEach(async () => { + jest.restoreAllMocks(); + await server.closeAndWait(); + }); + + it('connects a streaming request when the CA is supplied in the TLS options', async () => { + const requests = new NodeRequests({ ca: server.certificate }); + const res = await requests.fetch(server.url, { method: 'GET', streaming: true }); + expect(res.status).toEqual(200); + }); + + it('passes the CA to the per-request options of an https streaming request', async () => { + const requestSpy = jest.spyOn(https, 'request'); + + const requests = new NodeRequests({ ca: server.certificate }); + await requests.fetch(server.url, { method: 'GET', streaming: true }); + + expect(requestSpy).toHaveBeenCalledTimes(1); + expect(requestSpy.mock.calls[0][1]).toEqual( + expect.objectContaining({ ca: server.certificate }), + ); + }); + + it('rejects a streaming request when no CA is supplied', async () => { + const requests = new NodeRequests(); + await expect(requests.fetch(server.url, { method: 'GET', streaming: true })).rejects.toThrow(); + }); + + it('fails the stream against an untrusted certificate by default', async () => { + const requests = new NodeRequests(); + const es = requests.createEventSource(`${server.url}/stream`, { + headers: {}, + initialRetryDelayMillis: 100, + readTimeoutMillis: 5000, + retryResetIntervalMillis: 30_000, + errorFilter: () => false, + }); + try { + const errors = new AsyncQueue(); + es.addEventListener('error', (event) => errors.add(event)); + expect(await errors.take()).toBeDefined(); + } finally { + es.close(); + } + }); + + it('streams over TLS when the CA is supplied in the TLS options', async () => { + const chunks = new AsyncQueue(); + chunks.add('data: secure\n\n'); + server.byDefault( + TestHttpHandlers.chunkedStream(200, { 'content-type': 'text/event-stream' }, chunks), + ); + + const requests = new NodeRequests({ ca: server.certificate }); + const es = requests.createEventSource(`${server.url}/stream`, { + headers: {}, + initialRetryDelayMillis: 100, + readTimeoutMillis: 5000, + retryResetIntervalMillis: 30_000, + errorFilter: () => false, + }); + try { + const messages = new AsyncQueue<{ data?: string }>(); + es.addEventListener('message', (event) => messages.add(event ?? {})); + const message = await messages.take(); + expect(message.data).toEqual('secure'); + } finally { + es.close(); + } + }); + + it('does not pass the CA to the per-request options when a proxyAgent is supplied', async () => { + const requestSpy = jest.spyOn(https, 'request'); + const logger = { warn: jest.fn(), error: jest.fn(), info: jest.fn(), debug: jest.fn() }; + + const requests = new NodeRequests( + { ca: server.certificate }, + undefined, + new https.Agent({ ca: server.certificate }), + logger, + ); + await requests.fetch(server.url, { method: 'GET', streaming: true }); + + expect(requestSpy).toHaveBeenCalledTimes(1); + expect(requestSpy.mock.calls[0][1]).not.toHaveProperty('ca'); + }); +}); diff --git a/packages/sdk/server-node/package.json b/packages/sdk/server-node/package.json index c34f048029..c21ee19c29 100644 --- a/packages/sdk/server-node/package.json +++ b/packages/sdk/server-node/package.json @@ -47,9 +47,9 @@ }, "license": "Apache-2.0", "dependencies": { + "@launchdarkly/eventsource": "0.1.0", "@launchdarkly/js-server-sdk-common": "2.24.1", - "https-proxy-agent": "^7.0.6", - "launchdarkly-eventsource": "2.3.0" + "https-proxy-agent": "^7.0.6" }, "devDependencies": { "@types/jest": "^29.4.0", diff --git a/packages/sdk/server-node/src/platform/HeaderWrapper.ts b/packages/sdk/server-node/src/platform/HeaderWrapper.ts index dee99edc92..fe7781baf0 100644 --- a/packages/sdk/server-node/src/platform/HeaderWrapper.ts +++ b/packages/sdk/server-node/src/platform/HeaderWrapper.ts @@ -52,6 +52,17 @@ export default class HeaderWrapper implements platform.Headers { } } + /** + * Executes the callback once for each header, with the value first. The order matches the + * fetch `Headers.forEach` signature. Multi-value headers are joined with a comma, and + * headers without a value are skipped, like `entries`. + */ + forEach(callback: (value: string, key: string) => void): void { + for (const [key, value] of this.entries()) { + callback(value, key); + } + } + has(name: string): boolean { return Object.prototype.hasOwnProperty.call(this._headers, name); } diff --git a/packages/sdk/server-node/src/platform/NodeRequests.ts b/packages/sdk/server-node/src/platform/NodeRequests.ts index fddce38559..45cde40c41 100644 --- a/packages/sdk/server-node/src/platform/NodeRequests.ts +++ b/packages/sdk/server-node/src/platform/NodeRequests.ts @@ -1,15 +1,14 @@ import * as http from 'http'; import * as https from 'https'; import { HttpsProxyAgent, HttpsProxyAgentOptions } from 'https-proxy-agent'; -// No types for the event source. -// @ts-ignore -import { EventSource as LDEventSource } from 'launchdarkly-eventsource'; import { format as formatUrl } from 'url'; import { promisify } from 'util'; import * as zlib from 'zlib'; +import { createEventSource } from '@launchdarkly/eventsource'; import { EventSourceCapabilities, + internal, LDLogger, LDProxyOptions, LDTLSOptions, @@ -20,6 +19,37 @@ import NodeResponse from './NodeResponse'; const gzip = promisify(zlib.gzip); +/** + * The TLS options that are copied onto each streaming `https` request. The names match the + * options of `https.request()`. + */ +const TLS_OPTION_NAMES = [ + 'pfx', + 'key', + 'passphrase', + 'cert', + 'ca', + 'ciphers', + 'rejectUnauthorized', + 'secureProtocol', + 'servername', + 'checkServerIdentity', +] as const; + +function tlsRequestOptions(tlsOptions?: LDTLSOptions): Record { + const merged: Record = {}; + if (!tlsOptions) { + return merged; + } + const bag = tlsOptions as unknown as Record; + TLS_OPTION_NAMES.forEach((name) => { + if (bag[name] !== undefined) { + merged[name] = bag[name]; + } + }); + return merged; +} + function processTlsOptions(tlsOptions: LDTLSOptions): https.AgentOptions { const options: https.AgentOptions & { [index: string]: any } = { ca: tlsOptions.ca, @@ -122,7 +152,7 @@ function resolveAgent( export default class NodeRequests implements platform.Requests { private _agent: https.Agent | http.Agent | undefined; - private _tlsOptions: LDTLSOptions | undefined; + private _tlsParams: Record; private _hasProxy: boolean = false; @@ -138,6 +168,9 @@ export default class NodeRequests implements platform.Requests { enableEventCompression?: boolean, ) { this._agent = resolveAgent(tlsOptions, proxyOptions, proxyAgent, logger); + // The agent owns connection setup when the caller supplies one, so the per-request TLS + // parameters are only forwarded when this class built the agent itself (or no agent exists). + this._tlsParams = tlsRequestOptions(proxyAgent ? undefined : tlsOptions); // A caller-supplied proxyAgent is treated as a best-effort proxy signal: the SDK cannot // inspect an opaque agent to know whether it actually proxies (it could just as easily be a // certificate-only agent for mTLS). Reporting true is the better default here because @@ -153,6 +186,9 @@ export default class NodeRequests implements platform.Requests { } async fetch(url: string, options: platform.Options = {}): Promise { + if (options.streaming) { + return this._streamingFetch(url, options); + } const isSecure = url.startsWith('https://'); const impl = isSecure ? https : http; @@ -160,7 +196,7 @@ export default class NodeRequests implements platform.Requests { let bodyData: string | Buffer | undefined = options.body; // For get requests we are going to automatically support compressed responses. - // Note this does not affect SSE as the event source is not using this fetch implementation. + // Note this does not affect SSE as streaming requests take the branch above. if (options.method?.toLowerCase() === 'get') { headers['accept-encoding'] = 'gzip'; } @@ -205,18 +241,58 @@ export default class NodeRequests implements platform.Requests { }); } + /** + * The transport for a streaming request. It does not request compressed content, and it never + * follows a redirect. A redirect status resolves like any other non-200 response, and the + * caller decides whether to retry the original URL. It applies no read or socket timeout. The + * caller owns the read timeout and cancels through the abort signal. + */ + private _streamingFetch(url: string, options: platform.Options): Promise { + const isSecure = url.startsWith('https://'); + const impl = isSecure ? https : http; + const requestOptions: https.RequestOptions & Record = { + method: options.method, + headers: options.headers, + agent: this._agent, + }; + if (isSecure) { + Object.assign(requestOptions, this._tlsParams); + } + return new Promise((resolve, reject) => { + const req = impl.request(url, requestOptions, (res) => + resolve(internal.createStreamingResponse(res)), + ); + // An SSE consumer wants each chunk as soon as it arrives; do not batch small writes. + req.setNoDelay(true); + const { signal } = options; + if (signal) { + const abort = () => req.destroy(new Error('The stream request was aborted')); + if (signal.aborted) { + abort(); + } else { + signal.addEventListener('abort', abort, { once: true }); + } + } + // This listener stays attached after resolve. A later socket error then becomes a harmless + // no-op reject instead of an unhandled 'error' event that would crash the process + req.on('error', reject); + if (options.body !== undefined) { + req.write(options.body); + } + req.end(); + }); + } + createEventSource( url: string, eventSourceInitDict: platform.EventSourceInitDict, ): platform.EventSource { - const expandedOptions = { + return createEventSource(url, { ...eventSourceInitDict, - agent: this._agent, - tlsParams: this._tlsOptions, maxBackoffMillis: 30 * 1000, jitterRatio: 0.5, - }; - return new LDEventSource(url, expandedOptions); + fetch: (fetchUrl, init) => this.fetch(fetchUrl, { ...init, streaming: true }), + }); } getEventSourceCapabilities(): EventSourceCapabilities {