Skip to content
This repository was archived by the owner on Jan 23, 2026. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions .changeset/hip-actors-worry.md
Original file line number Diff line number Diff line change
@@ -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
6 changes: 1 addition & 5 deletions src/apis/api.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down
135 changes: 119 additions & 16 deletions src/apis/stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,18 +4,61 @@ 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
*/
class StreamImpl extends WritableStream implements Stream {
public readonly id: string;
public readonly url: string;
private activeWriter: WritableStreamDefaultWriter<Uint8Array> | 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<void> {
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);
}

/**
Expand All @@ -24,7 +67,15 @@ class StreamImpl extends WritableStream implements Stream {
*/
async close(): Promise<void> {
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) {
Expand All @@ -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
Expand Down Expand Up @@ -201,6 +254,7 @@ export default class StreamAPIImpl implements StreamAPI {
let putRequestPromise: Promise<Response> | 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
Expand All @@ -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<Uint8Array, Uint8Array>();

// Set up compression pipeline
const gzipStream = createGzip();
const nodeWritable = Writable.toWeb(
gzipStream
) as WritableStream<Uint8Array>;

// Pipe gzip output to the compressed readable
const gzipReader = Readable.toWeb(
gzipStream
) as ReadableStream<Uint8Array>;
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
Expand All @@ -226,14 +314,20 @@ export default class StreamAPIImpl implements StreamAPI {
}
const sdkVersion = getSDKVersion();

const headers: Record<string, string> = {
'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',
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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;
Expand Down
4 changes: 2 additions & 2 deletions src/router/context.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ export default class AgentContextWaitUntilHandler {
);
}
const currentContext = context.active();

// Start execution immediately, don't defer it
const executingPromise = (async () => {
running++;
Expand All @@ -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);
}
Expand Down
21 changes: 21 additions & 0 deletions src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Comment thread
jhaynie marked this conversation as resolved.
}

/**
Expand All @@ -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<void>;
/**
* close the stream gracefully, handling already closed streams without error
*/
Expand Down
Loading
Loading