diff --git a/.changeset/hip-actors-worry.md b/.changeset/hip-actors-worry.md new file mode 100644 index 00000000..b09ab001 --- /dev/null +++ b/.changeset/hip-actors-worry.md @@ -0,0 +1,8 @@ +--- +"@agentuity/sdk": patch +--- + +- Added support for automatic stream compression +- Added support for direct write to Stream in addition to getWriter() +- Added property `bytesWritten` to the Stream interface which represents the number of bytes written to the stream +- Added property `compressed` to the Stream interface which represents if the stream has compression enabled diff --git a/src/apis/api.ts b/src/apis/api.ts index 788bd597..89b852bc 100644 --- a/src/apis/api.ts +++ b/src/apis/api.ts @@ -28,11 +28,7 @@ interface ApiRequestWithUrl { type ApiRequestOptions = ApiRequestWithPath | ApiRequestWithUrl; -export type ServiceName = - | 'vector' - | 'keyvalue' - | 'stream' - | 'objectstore'; +export type ServiceName = 'vector' | 'keyvalue' | 'stream' | 'objectstore'; interface ApiRequestBase { method: 'POST' | 'GET' | 'PUT' | 'DELETE'; diff --git a/src/apis/stream.ts b/src/apis/stream.ts index a0d4a7dd..8bcd87e7 100644 --- a/src/apis/stream.ts +++ b/src/apis/stream.ts @@ -4,6 +4,7 @@ import type { CreateStreamProps, Stream, StreamAPI } from '../types'; import { context, SpanStatusCode, trace } from '@opentelemetry/api'; import { safeStringify } from '../utils/stringify'; import { ReadableStream } from 'node:stream/web'; +import { createGzip } from 'node:zlib'; /** * A writable stream implementation that extends WritableStream @@ -11,11 +12,53 @@ import { ReadableStream } from 'node:stream/web'; class StreamImpl extends WritableStream implements Stream { public readonly id: string; public readonly url: string; + private activeWriter: WritableStreamDefaultWriter | null = null; + public _bytesWritten = 0; + private _compressed: boolean; - constructor(id: string, url: string, underlyingSink: UnderlyingSink) { + constructor( + id: string, + url: string, + compressed: boolean, + underlyingSink: UnderlyingSink + ) { super(underlyingSink); this.id = id; this.url = url; + this._compressed = compressed; + } + + get bytesWritten(): number { + return this._bytesWritten; + } + + get compressed(): boolean { + return this._compressed; + } + + /** + * Write data to the stream + */ + async write( + chunk: string | Uint8Array | ArrayBuffer | Buffer | object + ): Promise { + let binaryChunk: Uint8Array; + if (chunk instanceof Uint8Array) { + binaryChunk = chunk; + } else if (typeof chunk === 'string') { + binaryChunk = new TextEncoder().encode(chunk); + } else if (chunk instanceof ArrayBuffer) { + binaryChunk = new Uint8Array(chunk); + } else if (typeof chunk === 'object' && chunk !== null) { + binaryChunk = new TextEncoder().encode(safeStringify(chunk)); + } else { + binaryChunk = new TextEncoder().encode(String(chunk)); + } + + if (!this.activeWriter) { + this.activeWriter = this.getWriter(); + } + await this.activeWriter.write(binaryChunk); } /** @@ -24,7 +67,15 @@ class StreamImpl extends WritableStream implements Stream { */ async close(): Promise { try { - // Check if stream is already closed by attempting to get a writer + // If we have an active writer from write() calls, use that + if (this.activeWriter) { + const writer = this.activeWriter; + this.activeWriter = null; + await writer.close(); + return; + } + + // Otherwise, get a writer and close it const writer = this.getWriter(); await writer.close(); } catch (error) { @@ -40,13 +91,15 @@ class StreamImpl extends WritableStream implements Stream { return Promise.resolve(); } // If the stream is locked, try to close the underlying writer - if ( - error instanceof TypeError && - error.message.includes('locked') - ) { + if (error instanceof TypeError && error.message.includes('locked')) { + // If we have an active writer, close it + if (this.activeWriter) { + const writer = this.activeWriter; + this.activeWriter = null; + await writer.close(); + return; + } // Best-effort closure for locked streams - // Note: We can't directly access the active writer, so we silently return - // In a real implementation, we would track the writer and close it here return Promise.resolve(); } // Re-throw any other errors @@ -201,6 +254,7 @@ export default class StreamAPIImpl implements StreamAPI { let putRequestPromise: Promise | null = null; let total = 0; let closed = false; + let streamInstance: StreamImpl | null = null; // Create a WritableStream that writes to the backend stream // Create the underlying sink that will handle the actual streaming @@ -210,10 +264,44 @@ export default class StreamAPIImpl implements StreamAPI { abortController = new AbortController(); // Create a ReadableStream to pipe data to the PUT request - const { readable, writable } = new TransformStream< + let { readable, writable } = new TransformStream< Uint8Array, Uint8Array >(); + + // If compression is enabled, add gzip transform + if (props?.compress) { + const { Readable, Writable } = await import('node:stream'); + + // Create a new transform for the compressed output + const { + readable: compressedReadable, + writable: compressedWritable, + } = new TransformStream(); + + // Set up compression pipeline + const gzipStream = createGzip(); + const nodeWritable = Writable.toWeb( + gzipStream + ) as WritableStream; + + // Pipe gzip output to the compressed readable + const gzipReader = Readable.toWeb( + gzipStream + ) as ReadableStream; + gzipReader.pipeTo(compressedWritable).catch((error) => { + abortController?.abort(error); + writer?.abort(error).catch(() => {}); + }); + + // Chain: writable -> gzip -> compressedReadable + readable.pipeTo(nodeWritable).catch((error) => { + abortController?.abort(error); + writer?.abort(error).catch(() => {}); + }); + readable = compressedReadable; + } + writer = writable.getWriter(); // Start the PUT request with the readable stream as body @@ -226,14 +314,20 @@ export default class StreamAPIImpl implements StreamAPI { } const sdkVersion = getSDKVersion(); + const headers: Record = { + 'Content-Type': + props?.contentType || 'application/octet-stream', + 'User-Agent': `Agentuity JS SDK/${sdkVersion}`, + Authorization: `Bearer ${apiKey}`, + }; + + if (props?.compress) { + headers['Content-Encoding'] = 'gzip'; + } + putRequestPromise = getFetch()(url, { method: 'PUT', - headers: { - 'Content-Type': - props?.contentType || 'application/octet-stream', - 'User-Agent': `Agentuity JS SDK/${sdkVersion}`, - Authorization: `Bearer ${apiKey}`, - }, + headers, body: readable, signal: abortController.signal, duplex: 'half', @@ -263,6 +357,9 @@ export default class StreamAPIImpl implements StreamAPI { // Write the chunk to the transform stream, which pipes to the PUT request await writer.write(binaryChunk); total += binaryChunk.length; + if (streamInstance) { + streamInstance._bytesWritten = total; + } }, async close() { if (closed) { @@ -306,7 +403,13 @@ export default class StreamAPIImpl implements StreamAPI { }, }; - const stream = new StreamImpl(result.id, url, underlyingSink); + const stream = new StreamImpl( + result.id, + url, + props?.compress ?? false, + underlyingSink + ); + streamInstance = stream; span.setStatus({ code: SpanStatusCode.OK }); return stream; diff --git a/src/router/context.ts b/src/router/context.ts index 1a9d32b2..c8dd8828 100644 --- a/src/router/context.ts +++ b/src/router/context.ts @@ -33,7 +33,7 @@ export default class AgentContextWaitUntilHandler { ); } const currentContext = context.active(); - + // Start execution immediately, don't defer it const executingPromise = (async () => { running++; @@ -58,7 +58,7 @@ export default class AgentContextWaitUntilHandler { } // NOTE: we only decrement when the promise is removed from the array in waitUntilAll })(); - + // Store the executing promise for cleanup tracking this.promises.push(executingPromise); } diff --git a/src/types.ts b/src/types.ts index 25cead85..0532561d 100644 --- a/src/types.ts +++ b/src/types.ts @@ -342,6 +342,13 @@ export interface CreateStreamProps { * optional contentType for the stream data. If not set, defaults to application/octet-stream */ contentType?: string; + + /** + * optional flag to enable gzip compression of stream data during upload. if true, will also add + * add Content-Encoding: gzip header to responses. The client MUST be able to accept gzip + * compression for this to work or must be able to uncompress the raw data it receives. + */ + compress?: true; } /** @@ -363,6 +370,20 @@ export interface Stream extends WritableStream { * the unique stream url to consume the stream */ url: string; + /** + * the total number of bytes written to the stream + */ + readonly bytesWritten: number; + /** + * whether the stream is using compression + */ + readonly compressed: boolean; + /** + * write data to the stream + */ + write( + chunk: string | Uint8Array | ArrayBuffer | Buffer | object + ): Promise; /** * close the stream gracefully, handling already closed streams without error */ diff --git a/test/apis/api.test.ts b/test/apis/api.test.ts index 14242415..a21a1cef 100644 --- a/test/apis/api.test.ts +++ b/test/apis/api.test.ts @@ -1,5 +1,14 @@ import { describe, expect, it, mock, beforeEach, afterEach } from 'bun:test'; -import { send, GET, POST, PUT, DELETE, setFetch, getFetch, getBaseUrlForService } from '../../src/apis/api'; +import { + send, + GET, + POST, + PUT, + DELETE, + setFetch, + getFetch, + getBaseUrlForService, +} from '../../src/apis/api'; import { createMockFetch } from '../setup'; import { ReadableStream } from 'node:stream/web'; @@ -298,7 +307,9 @@ describe('API Client', () => { expect(getBaseUrlForService()).toBe('https://agentuity.ai/'); expect(getBaseUrlForService('vector')).toBe('https://agentuity.ai/'); expect(getBaseUrlForService('keyvalue')).toBe('https://agentuity.ai/'); - expect(getBaseUrlForService('stream')).toBe('https://streams.agentuity.cloud'); + expect(getBaseUrlForService('stream')).toBe( + 'https://streams.agentuity.cloud' + ); expect(getBaseUrlForService('objectstore')).toBe('https://agentuity.ai/'); }); @@ -306,46 +317,71 @@ describe('API Client', () => { process.env.AGENTUITY_TRANSPORT_URL = 'https://transport.example.com/'; expect(getBaseUrlForService()).toBe('https://transport.example.com/'); - expect(getBaseUrlForService('vector')).toBe('https://transport.example.com/'); - expect(getBaseUrlForService('keyvalue')).toBe('https://transport.example.com/'); - expect(getBaseUrlForService('stream')).toBe('https://streams.agentuity.cloud'); // Stream service has its own default - expect(getBaseUrlForService('objectstore')).toBe('https://transport.example.com/'); + expect(getBaseUrlForService('vector')).toBe( + 'https://transport.example.com/' + ); + expect(getBaseUrlForService('keyvalue')).toBe( + 'https://transport.example.com/' + ); + expect(getBaseUrlForService('stream')).toBe( + 'https://streams.agentuity.cloud' + ); // Stream service has its own default + expect(getBaseUrlForService('objectstore')).toBe( + 'https://transport.example.com/' + ); }); it('should use service-specific URL for vector service', () => { process.env.AGENTUITY_TRANSPORT_URL = 'https://transport.example.com/'; process.env.AGENTUITY_VECTOR_URL = 'https://vector.example.com/'; - expect(getBaseUrlForService('vector')).toBe('https://vector.example.com/'); + expect(getBaseUrlForService('vector')).toBe( + 'https://vector.example.com/' + ); // Other services should still use transport URL - expect(getBaseUrlForService('keyvalue')).toBe('https://transport.example.com/'); + expect(getBaseUrlForService('keyvalue')).toBe( + 'https://transport.example.com/' + ); }); it('should use service-specific URL for keyvalue service', () => { process.env.AGENTUITY_TRANSPORT_URL = 'https://transport.example.com/'; process.env.AGENTUITY_KEYVALUE_URL = 'https://keyvalue.example.com/'; - expect(getBaseUrlForService('keyvalue')).toBe('https://keyvalue.example.com/'); + expect(getBaseUrlForService('keyvalue')).toBe( + 'https://keyvalue.example.com/' + ); // Other services should still use transport URL - expect(getBaseUrlForService('vector')).toBe('https://transport.example.com/'); + expect(getBaseUrlForService('vector')).toBe( + 'https://transport.example.com/' + ); }); it('should use service-specific URL for stream service', () => { process.env.AGENTUITY_TRANSPORT_URL = 'https://transport.example.com/'; process.env.AGENTUITY_STREAM_URL = 'https://stream.example.com/'; - expect(getBaseUrlForService('stream')).toBe('https://stream.example.com/'); + expect(getBaseUrlForService('stream')).toBe( + 'https://stream.example.com/' + ); // Other services should still use transport URL - expect(getBaseUrlForService('vector')).toBe('https://transport.example.com/'); + expect(getBaseUrlForService('vector')).toBe( + 'https://transport.example.com/' + ); }); it('should use service-specific URL for objectstore service', () => { process.env.AGENTUITY_TRANSPORT_URL = 'https://transport.example.com/'; - process.env.AGENTUITY_OBJECTSTORE_URL = 'https://objectstore.example.com/'; + process.env.AGENTUITY_OBJECTSTORE_URL = + 'https://objectstore.example.com/'; - expect(getBaseUrlForService('objectstore')).toBe('https://objectstore.example.com/'); + expect(getBaseUrlForService('objectstore')).toBe( + 'https://objectstore.example.com/' + ); // Other services should still use transport URL - expect(getBaseUrlForService('vector')).toBe('https://transport.example.com/'); + expect(getBaseUrlForService('vector')).toBe( + 'https://transport.example.com/' + ); }); it('should prioritize service-specific URLs over transport URL', () => { @@ -353,17 +389,28 @@ describe('API Client', () => { process.env.AGENTUITY_VECTOR_URL = 'https://vector.example.com/'; process.env.AGENTUITY_KEYVALUE_URL = 'https://keyvalue.example.com/'; process.env.AGENTUITY_STREAM_URL = 'https://stream.example.com/'; - process.env.AGENTUITY_OBJECTSTORE_URL = 'https://objectstore.example.com/'; + process.env.AGENTUITY_OBJECTSTORE_URL = + 'https://objectstore.example.com/'; - expect(getBaseUrlForService('vector')).toBe('https://vector.example.com/'); - expect(getBaseUrlForService('keyvalue')).toBe('https://keyvalue.example.com/'); - expect(getBaseUrlForService('stream')).toBe('https://stream.example.com/'); - expect(getBaseUrlForService('objectstore')).toBe('https://objectstore.example.com/'); + expect(getBaseUrlForService('vector')).toBe( + 'https://vector.example.com/' + ); + expect(getBaseUrlForService('keyvalue')).toBe( + 'https://keyvalue.example.com/' + ); + expect(getBaseUrlForService('stream')).toBe( + 'https://stream.example.com/' + ); + expect(getBaseUrlForService('objectstore')).toBe( + 'https://objectstore.example.com/' + ); }); it('should handle undefined service parameter', () => { process.env.AGENTUITY_TRANSPORT_URL = 'https://transport.example.com/'; - expect(getBaseUrlForService(undefined)).toBe('https://transport.example.com/'); + expect(getBaseUrlForService(undefined)).toBe( + 'https://transport.example.com/' + ); }); }); @@ -373,7 +420,8 @@ describe('API Client', () => { process.env.AGENTUITY_VECTOR_URL = 'https://vector.example.com/'; process.env.AGENTUITY_KEYVALUE_URL = 'https://keyvalue.example.com/'; process.env.AGENTUITY_STREAM_URL = 'https://stream.example.com/'; - process.env.AGENTUITY_OBJECTSTORE_URL = 'https://objectstore.example.com/'; + process.env.AGENTUITY_OBJECTSTORE_URL = + 'https://objectstore.example.com/'; }); it('should use vector service URL for GET request', async () => { @@ -385,7 +433,14 @@ describe('API Client', () => { }); it('should use keyvalue service URL for POST request', async () => { - await POST('/test', JSON.stringify({ data: 'test' }), undefined, undefined, undefined, 'keyvalue'); + await POST( + '/test', + JSON.stringify({ data: 'test' }), + undefined, + undefined, + undefined, + 'keyvalue' + ); expect(fetchCalls.length).toBeGreaterThan(0); const [url] = fetchCalls[0]; @@ -393,7 +448,14 @@ describe('API Client', () => { }); it('should use stream service URL for PUT request', async () => { - await PUT('/test', JSON.stringify({ data: 'test' }), undefined, undefined, undefined, 'stream'); + await PUT( + '/test', + JSON.stringify({ data: 'test' }), + undefined, + undefined, + undefined, + 'stream' + ); expect(fetchCalls.length).toBeGreaterThan(0); const [url] = fetchCalls[0]; @@ -401,7 +463,14 @@ describe('API Client', () => { }); it('should use objectstore service URL for DELETE request', async () => { - await DELETE('/test', undefined, undefined, undefined, undefined, 'objectstore'); + await DELETE( + '/test', + undefined, + undefined, + undefined, + undefined, + 'objectstore' + ); expect(fetchCalls.length).toBeGreaterThan(0); const [url] = fetchCalls[0]; @@ -427,60 +496,101 @@ describe('API Client', () => { }); it('should handle POST request with auth token and service', async () => { - await POST('/test', JSON.stringify({ data: 'test' }), undefined, undefined, 'custom-token', 'vector'); + await POST( + '/test', + JSON.stringify({ data: 'test' }), + undefined, + undefined, + 'custom-token', + 'vector' + ); expect(fetchCalls.length).toBeGreaterThan(0); const [url, options] = fetchCalls[0]; expect(url.toString()).toBe('https://vector.example.com/test'); - + const headers = options?.headers as Record; expect(headers?.Authorization).toBe('Bearer custom-token'); }); it('should handle GET request with auth token and service', async () => { - await GET('/test', false, undefined, undefined, 'get-custom-token', 'keyvalue'); + await GET( + '/test', + false, + undefined, + undefined, + 'get-custom-token', + 'keyvalue' + ); expect(fetchCalls.length).toBeGreaterThan(0); const [url, options] = fetchCalls[0]; expect(url.toString()).toBe('https://keyvalue.example.com/test'); - + const headers = options?.headers as Record; expect(headers?.Authorization).toBe('Bearer get-custom-token'); }); it('should handle PUT request with auth token and service', async () => { - await PUT('/test', JSON.stringify({ data: 'test' }), undefined, undefined, 'put-custom-token', 'stream'); + await PUT( + '/test', + JSON.stringify({ data: 'test' }), + undefined, + undefined, + 'put-custom-token', + 'stream' + ); expect(fetchCalls.length).toBeGreaterThan(0); const [url, options] = fetchCalls[0]; expect(url.toString()).toBe('https://stream.example.com/test'); - + const headers = options?.headers as Record; expect(headers?.Authorization).toBe('Bearer put-custom-token'); }); it('should handle DELETE request with auth token and service', async () => { - await DELETE('/test', undefined, undefined, undefined, 'delete-custom-token', 'objectstore'); + await DELETE( + '/test', + undefined, + undefined, + undefined, + 'delete-custom-token', + 'objectstore' + ); expect(fetchCalls.length).toBeGreaterThan(0); const [url, options] = fetchCalls[0]; expect(url.toString()).toBe('https://objectstore.example.com/test'); - + const headers = options?.headers as Record; expect(headers?.Authorization).toBe('Bearer delete-custom-token'); }); it('should validate all services work correctly', async () => { - const services: Array<{ name: 'vector' | 'keyvalue' | 'stream' | 'objectstore', expectedUrl: string }> = [ + const services: Array<{ + name: 'vector' | 'keyvalue' | 'stream' | 'objectstore'; + expectedUrl: string; + }> = [ { name: 'vector', expectedUrl: 'https://vector.example.com/test' }, { name: 'keyvalue', expectedUrl: 'https://keyvalue.example.com/test' }, { name: 'stream', expectedUrl: 'https://stream.example.com/test' }, - { name: 'objectstore', expectedUrl: 'https://objectstore.example.com/test' }, + { + name: 'objectstore', + expectedUrl: 'https://objectstore.example.com/test', + }, ]; for (const service of services) { fetchCalls.length = 0; // Clear previous calls - await GET('/test', false, undefined, undefined, undefined, service.name); + await GET( + '/test', + false, + undefined, + undefined, + undefined, + service.name + ); expect(fetchCalls.length).toBeGreaterThan(0); const [url] = fetchCalls[0]; @@ -517,15 +627,34 @@ describe('API Client', () => { delete process.env.AGENTUITY_OBJECTSTORE_URL; const testCases = [ - { service: 'vector' as const, expectedUrl: 'https://vector.agentuity.ai/test' }, - { service: 'keyvalue' as const, expectedUrl: 'https://transport.agentuity.ai/test' }, - { service: 'stream' as const, expectedUrl: 'https://stream.agentuity.ai/test' }, - { service: 'objectstore' as const, expectedUrl: 'https://transport.agentuity.ai/test' }, + { + service: 'vector' as const, + expectedUrl: 'https://vector.agentuity.ai/test', + }, + { + service: 'keyvalue' as const, + expectedUrl: 'https://transport.agentuity.ai/test', + }, + { + service: 'stream' as const, + expectedUrl: 'https://stream.agentuity.ai/test', + }, + { + service: 'objectstore' as const, + expectedUrl: 'https://transport.agentuity.ai/test', + }, ]; for (const testCase of testCases) { fetchCalls.length = 0; - await GET('/test', false, undefined, undefined, undefined, testCase.service); + await GET( + '/test', + false, + undefined, + undefined, + undefined, + testCase.service + ); expect(fetchCalls.length).toBeGreaterThan(0); const [url] = fetchCalls[0]; @@ -537,12 +666,12 @@ describe('API Client', () => { describe('getFetch', () => { it('should return the currently set fetch function', () => { const originalFetch = getFetch(); - + const mockFetch = mock(() => Promise.resolve({} as Response)); setFetch(mockFetch as unknown as typeof fetch); - + expect(getFetch()).toBe(mockFetch); - + // Restore original setFetch(originalFetch); }); diff --git a/test/stream.test.ts b/test/stream.test.ts index abed5c2e..1fb8f8da 100644 --- a/test/stream.test.ts +++ b/test/stream.test.ts @@ -39,27 +39,29 @@ describe('StreamAPI', () => { describe('create', () => { it('should validate stream name length', async () => { // Test empty name - await expect( - streamAPI.create('') - ).rejects.toThrow('Stream name must be between 1 and 254 characters'); + await expect(streamAPI.create('')).rejects.toThrow( + 'Stream name must be between 1 and 254 characters' + ); // Test too long name const longName = 'a'.repeat(255); - await expect( - streamAPI.create(longName) - ).rejects.toThrow('Stream name must be between 1 and 254 characters'); + await expect(streamAPI.create(longName)).rejects.toThrow( + 'Stream name must be between 1 and 254 characters' + ); }); it('should accept valid stream props', async () => { const props: CreateStreamProps = { - metadata: { customerId: 'customer-123', type: 'llm-response' } + metadata: { customerId: 'customer-123', type: 'llm-response' }, }; // Mock a simple successful response for validation tests - const mockFetch = mock(() => Promise.resolve({ - status: 200, - response: { status: 200, statusText: 'OK' } - })); + const mockFetch = mock(() => + Promise.resolve({ + status: 200, + response: { status: 200, statusText: 'OK' }, + }) + ); setFetch(mockFetch as unknown as typeof fetch); // This should not throw for valid props @@ -73,10 +75,12 @@ describe('StreamAPI', () => { it('should accept minimal stream props', async () => { // Mock a simple successful response for validation tests - const mockFetch = mock(() => Promise.resolve({ - status: 200, - response: { status: 200, statusText: 'OK' } - })); + const mockFetch = mock(() => + Promise.resolve({ + status: 200, + response: { status: 200, statusText: 'OK' }, + }) + ); setFetch(mockFetch as unknown as typeof fetch); try { @@ -94,10 +98,12 @@ describe('StreamAPI', () => { const maxName = 'a'.repeat(254); // Mock a simple successful response for validation tests - const mockFetch = mock(() => Promise.resolve({ - status: 200, - response: { status: 200, statusText: 'OK' } - })); + const mockFetch = mock(() => + Promise.resolve({ + status: 200, + response: { status: 200, statusText: 'OK' }, + }) + ); setFetch(mockFetch as unknown as typeof fetch); // These should not throw validation errors, only network-related errors @@ -135,90 +141,96 @@ describe('StreamAPI', () => { streamReadingComplete = false; streamReadingCompleteResolve = null; - const mockFetch = mock(async (url: URL | RequestInfo, options?: RequestInit) => { - fetchCalls.push([url, options]); + const mockFetch = mock( + async (url: URL | RequestInfo, options?: RequestInit) => { + fetchCalls.push([url, options]); - // Handle POST request to create stream - if (options?.method === 'POST') { - return { - status: 200, - response: { - json: () => Promise.resolve({ id: 'stream-123' }), + // Handle POST request to create stream + if (options?.method === 'POST') { + return { status: 200, - statusText: 'OK', - }, - json: () => Promise.resolve({ id: 'stream-123' }), - headers: new Headers({ 'content-type': 'application/json' }), - }; - } - - // Handle PUT request to upload stream data - if (options?.method === 'PUT') { - capturedAbortSignal = options.signal as AbortSignal; - - // Create a promise that we can resolve manually - putRequestPromise = new Promise((resolve, reject) => { - putRequestResolve = resolve; - _putRequestReject = reject; - - // Handle abort signal - if (capturedAbortSignal) { - capturedAbortSignal.addEventListener('abort', () => { - reject(new DOMException('The operation was aborted.', 'AbortError')); - }); - } - }); - - // Capture the stream data if present - if (options.body instanceof ReadableStream) { - const reader = options.body.getReader(); + response: { + json: () => Promise.resolve({ id: 'stream-123' }), + status: 200, + statusText: 'OK', + }, + json: () => Promise.resolve({ id: 'stream-123' }), + headers: new Headers({ 'content-type': 'application/json' }), + }; + } - // Create a promise we can await for stream reading completion - const _streamReadingCompletePromise = new Promise((resolve) => { - streamReadingCompleteResolve = resolve; + // Handle PUT request to upload stream data + if (options?.method === 'PUT') { + capturedAbortSignal = options.signal as AbortSignal; + + // Create a promise that we can resolve manually + putRequestPromise = new Promise((resolve, reject) => { + putRequestResolve = resolve; + _putRequestReject = reject; + + // Handle abort signal + if (capturedAbortSignal) { + capturedAbortSignal.addEventListener('abort', () => { + reject( + new DOMException('The operation was aborted.', 'AbortError') + ); + }); + } }); - // Read all chunks from the stream in background - const readStream = async () => { - try { - while (true) { - const { done, value } = await reader.read(); - if (done) { - streamReadingComplete = true; - if (streamReadingCompleteResolve) { - streamReadingCompleteResolve(); + // Capture the stream data if present + if (options.body instanceof ReadableStream) { + const reader = options.body.getReader(); + + // Create a promise we can await for stream reading completion + const _streamReadingCompletePromise = new Promise( + (resolve) => { + streamReadingCompleteResolve = resolve; + } + ); + + // Read all chunks from the stream in background + const readStream = async () => { + try { + while (true) { + const { done, value } = await reader.read(); + if (done) { + streamReadingComplete = true; + if (streamReadingCompleteResolve) { + streamReadingCompleteResolve(); + } + break; } - break; + capturedStreamData.push(value); + } + } catch (_error) { + streamReadingComplete = true; + if (streamReadingCompleteResolve) { + streamReadingCompleteResolve(); } - capturedStreamData.push(value); } - } catch (_error) { + }; + + readStream().catch(() => { streamReadingComplete = true; if (streamReadingCompleteResolve) { streamReadingCompleteResolve(); } - } - }; + }); + } - readStream().catch(() => { - streamReadingComplete = true; - if (streamReadingCompleteResolve) { - streamReadingCompleteResolve(); - } - }); + return putRequestPromise; } - return putRequestPromise; - } - - return { - status: 404, - response: { + return { status: 404, - statusText: 'Not Found', - }, - }; - }); + response: { + status: 404, + statusText: 'Not Found', + }, + }; + } + ); setFetch(mockFetch as unknown as typeof fetch); // Also set globalThis.fetch for the internal fetch call in stream.ts @@ -299,7 +311,9 @@ describe('StreamAPI', () => { }, 10); // Closing should throw an error due to failed PUT request - await expect(writer.close()).rejects.toThrow('PUT request failed: 500 Internal Server Error'); + await expect(writer.close()).rejects.toThrow( + 'PUT request failed: 500 Internal Server Error' + ); }); it('should handle large data streams', async () => { @@ -354,7 +368,9 @@ describe('StreamAPI', () => { await writer.close(); // Writing after close should throw - await expect(writer.write(new TextEncoder().encode('After close'))).rejects.toThrow(); + await expect( + writer.write(new TextEncoder().encode('After close')) + ).rejects.toThrow(); }); it('should handle concurrent writes', async () => { @@ -449,7 +465,7 @@ describe('StreamAPI', () => { 'World! ', 'This is a test of data integrity. ', '🚀 Unicode works too! ', - 'Final chunk.' + 'Final chunk.', ]; const sourceStream = new ReadableStream({ @@ -459,7 +475,7 @@ describe('StreamAPI', () => { controller.enqueue(new TextEncoder().encode(data)); } controller.close(); - } + }, }); const stream = await streamAPI.create('test-stream'); @@ -480,7 +496,7 @@ describe('StreamAPI', () => { // Wait for stream reading to complete if (!streamReadingComplete) { - await new Promise(resolve => setTimeout(resolve, 100)); + await new Promise((resolve) => setTimeout(resolve, 100)); } // Verify that the PUT request was made with a ReadableStream body @@ -502,83 +518,88 @@ describe('StreamAPI', () => { it('should correctly convert different data types to Uint8Array', async () => { // Test the data type conversion logic directly by using a simpler mock const writtenChunks: Uint8Array[] = []; - - const simpleMockFetch = mock(async (_url: string | URL | Request, options?: RequestInit) => { - if (options?.method === 'POST') { - return { - status: 200, - json: () => Promise.resolve({ id: 'test-123' }), - headers: new Headers({ 'content-type': 'application/json' }), - }; - } - - if (options?.method === 'PUT') { - // Capture the actual written data for validation - if (options.body instanceof ReadableStream) { - const reader = options.body.getReader(); - const captureData = async () => { - try { - while (true) { - const { done, value } = await reader.read(); - if (done) break; - writtenChunks.push(value); - } - } catch (_error) { - // Ignore stream errors - } + + const simpleMockFetch = mock( + async (_url: string | URL | Request, options?: RequestInit) => { + if (options?.method === 'POST') { + return { + status: 200, + json: () => Promise.resolve({ id: 'test-123' }), + headers: new Headers({ 'content-type': 'application/json' }), }; - captureData().catch(() => {}); } - - return Promise.resolve({ - ok: true, - status: 200, - statusText: 'OK', - } as Response); + + if (options?.method === 'PUT') { + // Capture the actual written data for validation + if (options.body instanceof ReadableStream) { + const reader = options.body.getReader(); + const captureData = async () => { + try { + while (true) { + const { done, value } = await reader.read(); + if (done) break; + writtenChunks.push(value); + } + } catch (_error) { + // Ignore stream errors + } + }; + captureData().catch(() => {}); + } + + return Promise.resolve({ + ok: true, + status: 200, + statusText: 'OK', + } as Response); + } + + return { status: 404 }; } - - return { status: 404 }; - }); + ); setFetch(simpleMockFetch as unknown as typeof fetch); globalThis.fetch = simpleMockFetch as unknown as typeof fetch; - + const stream = await streamAPI.create('type-test-stream'); const writer = stream.getWriter(); - + // Test different data types await writer.write(new Uint8Array([65])); // 'A' as Uint8Array - await writer.write('B'); // 'B' as string - + await writer.write('B'); // 'B' as string + const arrayBuffer = new ArrayBuffer(1); new Uint8Array(arrayBuffer)[0] = 67; // 'C' - await writer.write(arrayBuffer); // 'C' as ArrayBuffer - + await writer.write(arrayBuffer); // 'C' as ArrayBuffer + // Test Node.js Buffer (should be handled as Uint8Array since Buffer extends Uint8Array) const buffer = Buffer.from('D'); - await writer.write(buffer); // 'D' as Buffer - + await writer.write(buffer); // 'D' as Buffer + // Test object (should be converted to JSON string) const testObject = { message: 'E', number: 42 }; - await writer.write(testObject); // Object as JSON - + await writer.write(testObject); // Object as JSON + await writer.close(); - + // Wait for data capture - await new Promise(resolve => setTimeout(resolve, 50)); - + await new Promise((resolve) => setTimeout(resolve, 50)); + // Verify the data was converted correctly expect(writtenChunks.length).toBeGreaterThan(0); - + // Combine all chunks to verify the final result - const totalLength = writtenChunks.reduce((sum, chunk) => sum + chunk.length, 0); + const totalLength = writtenChunks.reduce( + (sum, chunk) => sum + chunk.length, + 0 + ); const combined = new Uint8Array(totalLength); let offset = 0; for (const chunk of writtenChunks) { combined.set(chunk, offset); offset += chunk.length; } - + // Should spell "ABCD" + JSON object const result = new TextDecoder().decode(combined); expect(result).toBe('ABCD{"message":"E","number":42}'); @@ -587,16 +608,16 @@ describe('StreamAPI', () => { it('should handle object serialization correctly', async () => { const stream = await streamAPI.create('test-stream'); const writer = stream.getWriter(); - + // Test various object types await writer.write({ simple: 'object' }); await writer.write({ nested: { data: 'value' }, array: [1, 2, 3] }); await writer.write({ number: 42, boolean: true, null: null }); - + // Test edge cases await writer.write({}); // Empty object await writer.write([]); // Empty array - + // Close the writer setTimeout(() => { if (putRequestResolve) { @@ -607,7 +628,7 @@ describe('StreamAPI', () => { } as Response); } }, 10); - + await writer.close(); // Verify PUT request was made with ReadableStream body @@ -620,22 +641,22 @@ describe('StreamAPI', () => { it('should handle various data sizes and content', async () => { const stream = await streamAPI.create('test-stream'); const writer = stream.getWriter(); - + // Test different content types and sizes await writer.write(new Uint8Array([72, 101, 108, 108, 111])); // "Hello" as binary - await writer.write(' World'); // String - await writer.write(' 🚀 Unicode test! 🎉'); // Unicode string - await writer.write(''); // Empty string - await writer.write('x'.repeat(1000)); // Large string - + await writer.write(' World'); // String + await writer.write(' 🚀 Unicode test! 🎉'); // Unicode string + await writer.write(''); // Empty string + await writer.write('x'.repeat(1000)); // Large string + // Test ArrayBuffer const arrayBuffer = new ArrayBuffer(3); const view = new Uint8Array(arrayBuffer); view[0] = 33; // ! - view[1] = 32; // space + view[1] = 32; // space view[2] = 65; // A await writer.write(arrayBuffer); - + // Close the writer setTimeout(() => { if (putRequestResolve) { @@ -646,7 +667,7 @@ describe('StreamAPI', () => { } as Response); } }, 10); - + await writer.close(); // Verify PUT request was made with ReadableStream body @@ -660,7 +681,7 @@ describe('StreamAPI', () => { // Test with custom content type const stream = await streamAPI.create('json-stream', { contentType: 'application/json', - metadata: { type: 'json-data' } + metadata: { type: 'json-data' }, }); const writer = stream.getWriter(); @@ -727,7 +748,7 @@ describe('StreamAPI', () => { fetchCalls.length = 0; const stream = await streamAPI.create(testCase.name, { - contentType: testCase.contentType + contentType: testCase.contentType, }); const writer = stream.getWriter(); @@ -791,50 +812,52 @@ describe('StreamAPI', () => { fetchCalls = []; originalFetch = globalThis.fetch; - const mockFetch = mock(async (url: URL | RequestInfo, options?: RequestInit) => { - fetchCalls.push([url, options]); + const mockFetch = mock( + async (url: URL | RequestInfo, options?: RequestInit) => { + fetchCalls.push([url, options]); - // Handle POST request to create stream - if (options?.method === 'POST') { - return { - status: 200, - response: { + // Handle POST request to create stream + if (options?.method === 'POST') { + return { + status: 200, + response: { + json: () => Promise.resolve({ id: 'stream-123' }), + status: 200, + statusText: 'OK', + }, json: () => Promise.resolve({ id: 'stream-123' }), + headers: new Headers({ 'content-type': 'application/json' }), + }; + } + + // Handle GET request to read stream data + if (options?.method === 'GET') { + const testData = 'Hello World from stream!'; + const encoder = new TextEncoder(); + const chunks = encoder.encode(testData); + + return { + ok: true, status: 200, statusText: 'OK', - }, - json: () => Promise.resolve({ id: 'stream-123' }), - headers: new Headers({ 'content-type': 'application/json' }), - }; - } - - // Handle GET request to read stream data - if (options?.method === 'GET') { - const testData = 'Hello World from stream!'; - const encoder = new TextEncoder(); - const chunks = encoder.encode(testData); + body: new ReadableStream({ + start(controller) { + controller.enqueue(chunks); + controller.close(); + }, + }), + }; + } return { - ok: true, - status: 200, - statusText: 'OK', - body: new ReadableStream({ - start(controller) { - controller.enqueue(chunks); - controller.close(); - } - }), + status: 404, + response: { + status: 404, + statusText: 'Not Found', + }, }; } - - return { - status: 404, - response: { - status: 404, - statusText: 'Not Found', - }, - }; - }); + ); setFetch(mockFetch as unknown as typeof fetch); globalThis.fetch = mockFetch as unknown as typeof fetch; @@ -846,10 +869,10 @@ describe('StreamAPI', () => { it('should return a ReadableStream when calling getReader()', async () => { const stream = await streamAPI.create('test-stream'); - + // Verify getReader method exists and returns a ReadableStream expect(typeof stream.getReader).toBe('function'); - + const reader = stream.getReader(); expect(reader).toBeInstanceOf(ReadableStream); }); @@ -857,50 +880,56 @@ describe('StreamAPI', () => { it('should make a GET request to the stream URL when reading', async () => { const stream = await streamAPI.create('test-stream'); const readableStream = stream.getReader(); - + // Read from the stream to trigger the GET request const reader = readableStream.getReader(); const { value, done } = await reader.read(); - + expect(done).toBe(false); expect(value).toBeInstanceOf(Uint8Array); - + // Verify GET request was made to correct URL - const getRequest = fetchCalls.find(([, options]) => options?.method === 'GET'); + const getRequest = fetchCalls.find( + ([, options]) => options?.method === 'GET' + ); expect(getRequest).toBeDefined(); - expect(getRequest?.[0].toString()).toBe('https://stream.test.com/stream-123'); + expect(getRequest?.[0].toString()).toBe( + 'https://stream.test.com/stream-123' + ); }); it('should include proper headers in GET request', async () => { const stream = await streamAPI.create('test-stream'); const readableStream = stream.getReader(); - + // Read from the stream const reader = readableStream.getReader(); await reader.read(); - + // Verify GET request headers - const getRequest = fetchCalls.find(([, options]) => options?.method === 'GET'); + const getRequest = fetchCalls.find( + ([, options]) => options?.method === 'GET' + ); expect(getRequest?.[1]?.headers).toMatchObject({ 'User-Agent': 'Agentuity JS SDK/1.0.0', - 'Authorization': 'Bearer test-api-key', + Authorization: 'Bearer test-api-key', }); }); it('should stream data correctly from the GET response', async () => { const stream = await streamAPI.create('test-stream'); const readableStream = stream.getReader(); - + // Read all data from the stream const reader = readableStream.getReader(); const chunks: Uint8Array[] = []; - + while (true) { const { value, done } = await reader.read(); if (done) break; chunks.push(value); } - + // Combine chunks and verify data const totalLength = chunks.reduce((sum, chunk) => sum + chunk.length, 0); const combined = new Uint8Array(totalLength); @@ -909,7 +938,7 @@ describe('StreamAPI', () => { combined.set(chunk, offset); offset += chunk.length; } - + const result = new TextDecoder().decode(combined); expect(result).toBe('Hello World from stream!'); }); @@ -917,94 +946,102 @@ describe('StreamAPI', () => { it('should handle API key authentication errors', async () => { // First create the stream with API key const stream = await streamAPI.create('test-stream'); - + // Then temporarily remove API key before calling getReader const originalApiKey = process.env.AGENTUITY_API_KEY; delete process.env.AGENTUITY_API_KEY; delete process.env.AGENTUITY_SDK_KEY; - + const readableStream = stream.getReader(); - + // Reading should fail with missing API key error const reader = readableStream.getReader(); - await expect(reader.read()).rejects.toThrow('AGENTUITY_API_KEY or AGENTUITY_SDK_KEY is not set'); - + await expect(reader.read()).rejects.toThrow( + 'AGENTUITY_API_KEY or AGENTUITY_SDK_KEY is not set' + ); + // Restore API key process.env.AGENTUITY_API_KEY = originalApiKey; }); it('should handle HTTP error responses', async () => { // Mock a failing GET request - const errorMockFetch = mock(async (_url: URL | RequestInfo, options?: RequestInit) => { - if (options?.method === 'POST') { - return { - status: 200, - response: { - json: () => Promise.resolve({ id: 'stream-123' }), + const errorMockFetch = mock( + async (_url: URL | RequestInfo, options?: RequestInit) => { + if (options?.method === 'POST') { + return { status: 200, - statusText: 'OK', - }, - json: () => Promise.resolve({ id: 'stream-123' }), - headers: new Headers({ 'content-type': 'application/json' }), - }; - } + response: { + json: () => Promise.resolve({ id: 'stream-123' }), + status: 200, + statusText: 'OK', + }, + json: () => Promise.resolve({ id: 'stream-123' }), + headers: new Headers({ 'content-type': 'application/json' }), + }; + } - if (options?.method === 'GET') { - return { - ok: false, - status: 500, - statusText: 'Internal Server Error', - }; - } + if (options?.method === 'GET') { + return { + ok: false, + status: 500, + statusText: 'Internal Server Error', + }; + } - return { status: 404 }; - }); + return { status: 404 }; + } + ); setFetch(errorMockFetch as unknown as typeof fetch); globalThis.fetch = errorMockFetch as unknown as typeof fetch; const stream = await streamAPI.create('test-stream'); const readableStream = stream.getReader(); - + // Reading should fail with HTTP error const reader = readableStream.getReader(); - await expect(reader.read()).rejects.toThrow('Failed to read stream: 500 Internal Server Error'); + await expect(reader.read()).rejects.toThrow( + 'Failed to read stream: 500 Internal Server Error' + ); }); it('should handle null response body', async () => { // Mock a GET request with null body - const nullBodyMockFetch = mock(async (_url: URL | RequestInfo, options?: RequestInit) => { - if (options?.method === 'POST') { - return { - status: 200, - response: { + const nullBodyMockFetch = mock( + async (_url: URL | RequestInfo, options?: RequestInit) => { + if (options?.method === 'POST') { + return { + status: 200, + response: { + json: () => Promise.resolve({ id: 'stream-123' }), + status: 200, + statusText: 'OK', + }, json: () => Promise.resolve({ id: 'stream-123' }), + headers: new Headers({ 'content-type': 'application/json' }), + }; + } + + if (options?.method === 'GET') { + return { + ok: true, status: 200, statusText: 'OK', - }, - json: () => Promise.resolve({ id: 'stream-123' }), - headers: new Headers({ 'content-type': 'application/json' }), - }; - } + body: null, // Null body + }; + } - if (options?.method === 'GET') { - return { - ok: true, - status: 200, - statusText: 'OK', - body: null, // Null body - }; + return { status: 404 }; } - - return { status: 404 }; - }); + ); setFetch(nullBodyMockFetch as unknown as typeof fetch); globalThis.fetch = nullBodyMockFetch as unknown as typeof fetch; const stream = await streamAPI.create('test-stream'); const readableStream = stream.getReader(); - + // Reading should fail with null body error const reader = readableStream.getReader(); await expect(reader.read()).rejects.toThrow('Response body is null'); @@ -1012,74 +1049,79 @@ describe('StreamAPI', () => { it('should handle large streams with multiple chunks', async () => { // Mock a GET request that returns multiple chunks - const multiChunkMockFetch = mock(async (_url: URL | RequestInfo, options?: RequestInit) => { - if (options?.method === 'POST') { - return { - status: 200, - response: { + const multiChunkMockFetch = mock( + async (_url: URL | RequestInfo, options?: RequestInit) => { + if (options?.method === 'POST') { + return { + status: 200, + response: { + json: () => Promise.resolve({ id: 'stream-123' }), + status: 200, + statusText: 'OK', + }, json: () => Promise.resolve({ id: 'stream-123' }), + headers: new Headers({ 'content-type': 'application/json' }), + }; + } + + if (options?.method === 'GET') { + const chunks = [ + 'Chunk 1: ', + 'Chunk 2: ', + 'Chunk 3: ', + 'Final chunk!', + ]; + + return { + ok: true, status: 200, statusText: 'OK', - }, - json: () => Promise.resolve({ id: 'stream-123' }), - headers: new Headers({ 'content-type': 'application/json' }), - }; - } - - if (options?.method === 'GET') { - const chunks = [ - 'Chunk 1: ', - 'Chunk 2: ', - 'Chunk 3: ', - 'Final chunk!' - ]; + body: new ReadableStream({ + start(controller) { + for (const chunk of chunks) { + controller.enqueue(new TextEncoder().encode(chunk)); + } + controller.close(); + }, + }), + }; + } - return { - ok: true, - status: 200, - statusText: 'OK', - body: new ReadableStream({ - start(controller) { - for (const chunk of chunks) { - controller.enqueue(new TextEncoder().encode(chunk)); - } - controller.close(); - } - }), - }; + return { status: 404 }; } - - return { status: 404 }; - }); + ); setFetch(multiChunkMockFetch as unknown as typeof fetch); globalThis.fetch = multiChunkMockFetch as unknown as typeof fetch; const stream = await streamAPI.create('test-stream'); const readableStream = stream.getReader(); - + // Read all chunks const reader = readableStream.getReader(); const receivedChunks: Uint8Array[] = []; - + while (true) { const { value, done } = await reader.read(); if (done) break; receivedChunks.push(value); } - + // Verify we received multiple chunks expect(receivedChunks.length).toBeGreaterThan(1); - + // Combine and verify final result - const totalLength = receivedChunks.reduce((sum, chunk) => sum + chunk.length, 0); + const totalLength = receivedChunks.reduce( + (sum, chunk) => sum + chunk.length, + 0 + ); const combined = new Uint8Array(totalLength); let offset = 0; for (const chunk of receivedChunks) { combined.set(chunk, offset); offset += chunk.length; } - + const result = new TextDecoder().decode(combined); expect(result).toBe('Chunk 1: Chunk 2: Chunk 3: Final chunk!'); }); @@ -1089,20 +1131,22 @@ describe('StreamAPI', () => { const originalApiKey = process.env.AGENTUITY_API_KEY; delete process.env.AGENTUITY_API_KEY; process.env.AGENTUITY_SDK_KEY = 'test-sdk-key'; - + const stream = await streamAPI.create('test-stream'); const readableStream = stream.getReader(); - + // Read from the stream const reader = readableStream.getReader(); await reader.read(); - + // Verify GET request used SDK key - const getRequest = fetchCalls.find(([, options]) => options?.method === 'GET'); + const getRequest = fetchCalls.find( + ([, options]) => options?.method === 'GET' + ); expect(getRequest?.[1]?.headers).toMatchObject({ - 'Authorization': 'Bearer test-sdk-key', + Authorization: 'Bearer test-sdk-key', }); - + // Restore original environment process.env.AGENTUITY_API_KEY = originalApiKey; delete process.env.AGENTUITY_SDK_KEY; @@ -1112,16 +1156,502 @@ describe('StreamAPI', () => { const stream = await streamAPI.create('test-stream'); const readableStream = stream.getReader(); const reader = readableStream.getReader(); - + // Read from the stream - should work immediately const { value, done } = await reader.read(); - + expect(done).toBe(false); expect(value).toBeInstanceOf(Uint8Array); - + // Verify GET request was made - const getRequest = fetchCalls.find(([, options]) => options?.method === 'GET'); + const getRequest = fetchCalls.find( + ([, options]) => options?.method === 'GET' + ); expect(getRequest).toBeDefined(); }); }); + + describe('direct write() and close() methods', () => { + let fetchCalls: Array<[URL | RequestInfo, RequestInit | undefined]>; + let putRequestPromise: Promise | null = null; + let putRequestResolve: ((response: Response) => void) | null = null; + + beforeEach(() => { + fetchCalls = []; + putRequestPromise = null; + putRequestResolve = null; + + const mockFetch = mock( + async (url: URL | RequestInfo, options?: RequestInit) => { + fetchCalls.push([url, options]); + + if (options?.method === 'POST') { + return { + status: 200, + response: { + json: () => Promise.resolve({ id: 'stream-123' }), + status: 200, + statusText: 'OK', + }, + json: () => Promise.resolve({ id: 'stream-123' }), + headers: new Headers({ 'content-type': 'application/json' }), + }; + } + + if (options?.method === 'PUT') { + // Drain the ReadableStream to prevent backpressure blocking + if (options.body instanceof ReadableStream) { + const reader = options.body.getReader(); + (async () => { + try { + while (true) { + const { done } = await reader.read(); + if (done) break; + } + } catch (_e) { + // Ignore errors from cancelled reads + } + })(); + } + + putRequestPromise = new Promise((resolve) => { + putRequestResolve = resolve; + }); + + return putRequestPromise; + } + + return { status: 404 }; + } + ); + + setFetch(mockFetch as unknown as typeof fetch); + globalThis.fetch = mockFetch as unknown as typeof fetch; + }); + + afterEach(() => { + globalThis.fetch = originalFetch; + }); + + it('should write directly to stream using write() method', async () => { + const stream = await streamAPI.create('test-stream'); + + await stream.write('Hello '); + await stream.write('World!'); + + // Verify PUT was initiated + expect(putRequestResolve).not.toBeNull(); + + // Simulate successful PUT response asynchronously (same pattern as getWriter tests) + setTimeout(() => { + if (putRequestResolve) { + putRequestResolve({ + ok: true, + status: 200, + statusText: 'OK', + } as Response); + } + }, 10); + + await stream.close(); + + expect(fetchCalls).toHaveLength(2); + const [, uploadOptions] = fetchCalls[1]; + expect(uploadOptions?.method).toBe('PUT'); + }); + + it('should write different data types using write() method', async () => { + const stream = await streamAPI.create('test-stream'); + + await stream.write('string'); + await stream.write(new TextEncoder().encode('Uint8Array')); + await stream.write(new ArrayBuffer(8)); + await stream.write({ key: 'value' }); + + setTimeout(() => { + if (putRequestResolve) { + putRequestResolve({ + ok: true, + status: 200, + statusText: 'OK', + } as Response); + } + }, 10); + + await stream.close(); + + expect(fetchCalls).toHaveLength(2); + }); + + it('should write using direct write() method multiple times', async () => { + const stream = await streamAPI.create('test-stream'); + + await stream.write('Direct write 1'); + await stream.write('Direct write 2'); + await stream.write('Direct write 3'); + + setTimeout(() => { + if (putRequestResolve) { + putRequestResolve({ + ok: true, + status: 200, + statusText: 'OK', + } as Response); + } + }, 10); + + await stream.close(); + + expect(fetchCalls).toHaveLength(2); + }); + + it('should verify write() method exists on Stream interface', async () => { + const stream = await streamAPI.create('test-stream'); + + expect(typeof stream.write).toBe('function'); + expect(typeof stream.close).toBe('function'); + expect(typeof stream.getWriter).toBe('function'); + }); + + it('should verify PUT request is initiated on stream creation', async () => { + const stream = await streamAPI.create('test-stream'); + + // PUT is initiated in the underlying sink's start() method + expect(fetchCalls.length).toBe(2); + expect(fetchCalls[0][1]?.method).toBe('POST'); + expect(fetchCalls[1][1]?.method).toBe('PUT'); + + // Stream is ready to write + expect(typeof stream.write).toBe('function'); + }); + + it('should have bytesWritten property initialized to 0', async () => { + const stream = await streamAPI.create('test-stream'); + expect(stream.bytesWritten).toBe(0); + expect(typeof stream.bytesWritten).toBe('number'); + }); + + it('should track bytesWritten correctly as data is written with write()', async () => { + const stream = await streamAPI.create('test-stream'); + + expect(stream.bytesWritten).toBe(0); + + // Write first chunk - "Hello" = 5 bytes + await stream.write('Hello'); + expect(stream.bytesWritten).toBe(5); + + // Write second chunk - " World" = 6 bytes + await stream.write(' World'); + expect(stream.bytesWritten).toBe(11); + + // Write third chunk - "!" = 1 byte + await stream.write('!'); + expect(stream.bytesWritten).toBe(12); + + // Clean up + setTimeout(() => { + if (putRequestResolve) { + putRequestResolve({ + ok: true, + status: 200, + statusText: 'OK', + } as Response); + } + }, 10); + + await stream.close(); + }); + + it('should track bytesWritten correctly when using getWriter()', async () => { + const stream = await streamAPI.create('test-stream'); + + expect(stream.bytesWritten).toBe(0); + + const writer = stream.getWriter(); + + // Write first chunk - "Test" = 4 bytes + await writer.write(new TextEncoder().encode('Test')); + expect(stream.bytesWritten).toBe(4); + + // Write second chunk - " Data" = 5 bytes + await writer.write(new TextEncoder().encode(' Data')); + expect(stream.bytesWritten).toBe(9); + + // Clean up + setTimeout(() => { + if (putRequestResolve) { + putRequestResolve({ + ok: true, + status: 200, + statusText: 'OK', + } as Response); + } + }, 10); + + await writer.close(); + }); + + it('should track bytesWritten correctly with different data types', async () => { + const stream = await streamAPI.create('test-stream'); + + expect(stream.bytesWritten).toBe(0); + + // String - "abc" = 3 bytes + await stream.write('abc'); + expect(stream.bytesWritten).toBe(3); + + // Uint8Array - 4 bytes + await stream.write(new Uint8Array([1, 2, 3, 4])); + expect(stream.bytesWritten).toBe(7); + + // Object (JSON) - {"x":1} = 7 bytes + await stream.write({ x: 1 }); + expect(stream.bytesWritten).toBe(14); + + // ArrayBuffer - 3 bytes + await stream.write(new ArrayBuffer(3)); + expect(stream.bytesWritten).toBe(17); + + // Clean up + setTimeout(() => { + if (putRequestResolve) { + putRequestResolve({ + ok: true, + status: 200, + statusText: 'OK', + } as Response); + } + }, 10); + + await stream.close(); + }); + + it('should have compressed property false by default', async () => { + const stream = await streamAPI.create('test-stream'); + expect(stream.compressed).toBe(false); + }); + + it('should have compressed property true when compression enabled', async () => { + const stream = await streamAPI.create('test-stream', { compress: true }); + expect(stream.compressed).toBe(true); + }); + }); + + describe('compression', () => { + let fetchCalls: Array<[URL | RequestInfo, RequestInit | undefined]>; + let putRequestResolve: ((response: Response) => void) | null = null; + let _putRequestReject: ((error: Error) => void) | null = null; + + beforeEach(() => { + fetchCalls = []; + putRequestResolve = null; + _putRequestReject = null; + + const mockFetch = mock( + async (url: URL | RequestInfo, options?: RequestInit) => { + fetchCalls.push([url, options]); + + if (options?.method === 'POST') { + return { + status: 200, + response: { + json: () => Promise.resolve({ id: 'stream-123' }), + status: 200, + statusText: 'OK', + }, + json: () => Promise.resolve({ id: 'stream-123' }), + headers: new Headers({ 'content-type': 'application/json' }), + }; + } + + if (options?.method === 'PUT') { + if (options.body instanceof ReadableStream) { + const reader = options.body.getReader(); + reader.read().catch(() => {}); + } + + const putRequestPromise = new Promise( + (resolve, reject) => { + putRequestResolve = resolve; + _putRequestReject = reject; + } + ); + return putRequestPromise; + } + + return { status: 404 }; + } + ); + + setFetch(mockFetch as unknown as typeof fetch); + globalThis.fetch = mockFetch as unknown as typeof fetch; + }); + + it('should set Content-Encoding: gzip header when compress is true', async () => { + const stream = await streamAPI.create('compressed-stream', { + compress: true, + }); + + const writer = stream.getWriter(); + await writer.write('test data'); + + setTimeout(() => { + if (putRequestResolve) { + putRequestResolve({ + ok: true, + status: 200, + statusText: 'OK', + } as Response); + } + }, 10); + + await writer.close(); + + expect(fetchCalls).toHaveLength(2); + const [, uploadOptions] = fetchCalls[1]; + expect(uploadOptions?.headers).toMatchObject({ + 'Content-Encoding': 'gzip', + }); + }); + + it('should not set Content-Encoding header when compress is not specified', async () => { + const stream = await streamAPI.create('uncompressed-stream'); + + const writer = stream.getWriter(); + await writer.write('test data'); + + setTimeout(() => { + if (putRequestResolve) { + putRequestResolve({ + ok: true, + status: 200, + statusText: 'OK', + } as Response); + } + }, 10); + + await writer.close(); + + expect(fetchCalls).toHaveLength(2); + const [, uploadOptions] = fetchCalls[1]; + const headers = uploadOptions?.headers as Record; + expect(headers?.['Content-Encoding']).toBeUndefined(); + }); + + it('should compress data when compress is true', async () => { + let capturedData: Uint8Array[] = []; + + const mockFetchWithCapture = mock( + async (url: URL | RequestInfo, options?: RequestInit) => { + fetchCalls.push([url, options]); + + if (options?.method === 'POST') { + return { + status: 200, + response: { + json: () => Promise.resolve({ id: 'stream-123' }), + status: 200, + statusText: 'OK', + }, + json: () => Promise.resolve({ id: 'stream-123' }), + headers: new Headers({ 'content-type': 'application/json' }), + }; + } + + if ( + options?.method === 'PUT' && + options.body instanceof ReadableStream + ) { + const reader = options.body.getReader(); + capturedData = []; + + const readStream = async () => { + while (true) { + const { done, value } = await reader.read(); + if (done) break; + capturedData.push(value); + } + }; + + readStream().catch(() => {}); + + return new Promise((resolve) => { + putRequestResolve = resolve; + }); + } + + return { status: 404 }; + } + ); + + setFetch(mockFetchWithCapture as unknown as typeof fetch); + globalThis.fetch = mockFetchWithCapture as unknown as typeof fetch; + + const testData = 'x'.repeat(1000); + const stream = await streamAPI.create('compressed-stream', { + compress: true, + }); + + const writer = stream.getWriter(); + await writer.write(testData); + + await new Promise((resolve) => setTimeout(resolve, 50)); + + setTimeout(() => { + if (putRequestResolve) { + putRequestResolve({ + ok: true, + status: 200, + statusText: 'OK', + } as Response); + } + }, 10); + + await writer.close(); + + await new Promise((resolve) => setTimeout(resolve, 50)); + + const compressedSize = capturedData.reduce( + (sum, chunk) => sum + chunk.length, + 0 + ); + const originalSize = new TextEncoder().encode(testData).length; + + expect(compressedSize).toBeGreaterThan(0); + expect(compressedSize).toBeLessThan(originalSize); + }); + + it('should propagate compression errors to abort the stream', async () => { + const { createGzip } = await import('node:zlib'); + const originalCreateGzip = createGzip; + + let errorThrown = false; + + mock.module('node:zlib', () => ({ + createGzip: () => { + const gzip = originalCreateGzip(); + setTimeout(() => { + gzip.destroy(new Error('Compression failed')); + errorThrown = true; + }, 20); + return gzip; + }, + })); + + const stream = await streamAPI.create('error-stream', { + compress: true, + }); + + const writer = stream.getWriter(); + + try { + await writer.write('test data'); + await new Promise((resolve) => setTimeout(resolve, 100)); + } catch (_error) { + errorThrown = true; + } + + expect(errorThrown).toBe(true); + }); + }); }); diff --git a/test/utils/stringify.test.ts b/test/utils/stringify.test.ts index d37502bf..40bcd77a 100644 --- a/test/utils/stringify.test.ts +++ b/test/utils/stringify.test.ts @@ -12,17 +12,17 @@ describe('safeStringify', () => { it('should handle circular references', () => { const obj = { name: 'test' } as Record; obj.self = obj; - + const result = safeStringify(obj); expect(result).toBe('{"name":"test","self":"[Circular]"}'); }); it('should handle bigint values', () => { - const obj = { + const obj = { regularNumber: 123, - bigNumber: 9007199254740991n + bigNumber: 9007199254740991n, }; - + const result = safeStringify(obj); expect(result).toBe('{"regularNumber":123,"bigNumber":"9007199254740991"}'); }); @@ -33,14 +33,14 @@ describe('safeStringify', () => { id: 1n, timestamp: BigInt(Date.now()), nested: { - value: 42n - } - } + value: 42n, + }, + }, }; - + const result = safeStringify(obj); const parsed = JSON.parse(result); - + expect(typeof parsed.data.id).toBe('string'); expect(parsed.data.id).toBe('1'); expect(typeof parsed.data.timestamp).toBe('string'); @@ -49,19 +49,19 @@ describe('safeStringify', () => { }); it('should handle bigint with circular references', () => { - const obj = { + const obj = { id: 123n, - name: 'test' + name: 'test', } as Record; obj.self = obj; - + const result = safeStringify(obj); expect(result).toBe('{"id":"123","name":"test","self":"[Circular]"}'); }); it('should handle arrays with bigint values', () => { const arr = [1n, 2n, { value: 3n }]; - + const result = safeStringify(arr); expect(result).toBe('[\"1\",\"2\",{\"value\":\"3\"}]'); }); @@ -71,12 +71,12 @@ describe('safeStringify', () => { zero: 0n, negative: -123n, maxSafe: BigInt(Number.MAX_SAFE_INTEGER), - large: 123456789012345678901234567890n + large: 123456789012345678901234567890n, }; - + const result = safeStringify(obj); const parsed = JSON.parse(result); - + expect(parsed.zero).toBe('0'); expect(parsed.negative).toBe('-123'); expect(parsed.maxSafe).toBe('9007199254740991'); @@ -89,17 +89,19 @@ describe('safeStringify', () => { first: sharedObject, second: sharedObject, nested: { - third: sharedObject - } + third: sharedObject, + }, }; - + const result = safeStringify(obj); - const expected = '{"first":{"name":"shared","value":42},"second":{"name":"shared","value":42},"nested":{"third":{"name":"shared","value":42}}}'; + const expected = + '{"first":{"name":"shared","value":42},"second":{"name":"shared","value":42},"nested":{"third":{"name":"shared","value":42}}}'; expect(result).toBe(expected); - + // Verify that the shared object appears multiple times, not as [Circular] expect(result).not.toContain('[Circular]'); - const occurrences = (result.match(/{"name":"shared","value":42}/g) || []).length; + const occurrences = (result.match(/{"name":"shared","value":42}/g) || []) + .length; expect(occurrences).toBe(3); }); });