Skip to content
Draft
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
54 changes: 54 additions & 0 deletions yellowstone-grpc-client-nodejs/__tests__/subscribe.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1218,6 +1218,60 @@ describe("Client connection guard behavior", () => {
});

describe("ClientDuplexStream read and lifecycle behavior", () => {
test("nativeStats: forwards native queue counters", async () => {
const stats = {
queuedMessages: 2,
queuedBytes: 1024,
maxQueuedBytes: 2048,
updatesEnqueued: 5,
updatesDequeued: 3,
grpcUpdatesReceived: 5,
readCalls: 4,
};
const nativeStats = jest.fn().mockReturnValue(stats);
const stream = new ClientDuplexStream(
{
close: jest.fn(),
writeRaw: jest.fn(),
read: jest.fn(),
stats: nativeStats,
},
{ objectMode: true },
);

expect(stream.nativeStats()).toBe(stats);
expect(nativeStats).toHaveBeenCalledTimes(1);

await closeStreamAndWait(stream);
});

test("deshred nativeStats: forwards native queue counters", async () => {
const stats = {
queuedMessages: 1,
queuedBytes: 512,
maxQueuedBytes: 1024,
updatesEnqueued: 2,
updatesDequeued: 1,
grpcUpdatesReceived: 2,
readCalls: 1,
};
const nativeStats = jest.fn().mockReturnValue(stats);
const stream = new ClientDeshredDuplexStream(
{
close: jest.fn(),
writeRaw: jest.fn(),
read: jest.fn(),
stats: nativeStats,
},
{ objectMode: true },
);

expect(stream.nativeStats()).toBe(stats);
expect(nativeStats).toHaveBeenCalledTimes(1);

await closeStreamAndWait(stream);
});

function makeNativeUpdate(): Uint8Array {
return encodeNativeSubscribeUpdate({
filters: ["client"],
Expand Down
14 changes: 14 additions & 0 deletions yellowstone-grpc-client-nodejs/napi/index.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ export declare class DuplexStream {
* Retrieve one encoded `SubscribeUpdate` payload from the worker.
*/
read(): Promise<Buffer | undefined | null>
/** Return current native stream queue counters for diagnostics. */
stats(): DuplexStreamStats
/** Close the stream and reject future writes. */
close(): void
writeRaw(requestBytes: Buffer): void
Expand All @@ -46,6 +48,8 @@ export declare class DuplexStream {
export declare class DuplexStreamDeshred {
/** Retrieve one encoded `SubscribeUpdateDeshred` payload. */
read(): Promise<Buffer | undefined | null>
/** Return current native stream queue counters for diagnostics. */
stats(): DuplexStreamStats
close(): void
writeRaw(requestBytes: Buffer): void
}
Expand Down Expand Up @@ -97,6 +101,16 @@ export const AUTORECONNECT_FILTER_KEY: string

export declare function decodeTxError(err: Array<number>): string

export interface DuplexStreamStats {
queuedMessages: number
queuedBytes: number
maxQueuedBytes: number
updatesEnqueued: number
updatesDequeued: number
grpcUpdatesReceived: number
readCalls: number
}

export declare function encodeDeshredTx(data: Uint8Array, encoding: WasmUiTransactionEncoding): string

export declare function encodeTx(data: Uint8Array, encoding: WasmUiTransactionEncoding, maxSupportedTransactionVersion: number | undefined | null, showRewards: boolean): string
Expand Down
Loading