Skip to content
Open
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
66 changes: 66 additions & 0 deletions llm/token-meter/tests/token-meter-anthropic.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
*/

import Stripe from 'stripe';
import {Stream as AnthropicStream} from '@anthropic-ai/sdk/streaming';
import {createTokenMeter} from '../token-meter';
import type {MeterConfig} from '../types';

Expand Down Expand Up @@ -231,6 +232,71 @@ describe('TokenMeter - Anthropic Provider', () => {
});

describe('Messages - Streaming', () => {
it('does not consume the source before the returned stream is read', async () => {
const meter = createTokenMeter(TEST_API_KEY, config);
let pulls = 0;

async function* chunks() {
pulls += 1;
yield {
type: 'message_start',
message: {
id: 'msg_lazy',
model: 'claude-3-5-sonnet-20241022',
usage: {input_tokens: 10, output_tokens: 0},
},
};
}

const source = new AnthropicStream(chunks, new AbortController());
const wrapped = meter.trackUsageStreamAnthropic(source as any, 'cus_123');

await new Promise(resolve => setImmediate(resolve));

expect(pulls).toBe(0);
expect(wrapped).toBeInstanceOf(AnthropicStream);
});

it('closes the source when the returned stream is abandoned', async () => {
const meter = createTokenMeter(TEST_API_KEY, config);
let finalized = false;
let release!: () => void;
const blocked = new Promise<void>(resolve => {
release = resolve;
});

async function* chunks() {
try {
yield {
type: 'message_start',
message: {
id: 'msg_cancel',
model: 'claude-3-5-sonnet-20241022',
usage: {input_tokens: 10, output_tokens: 0},
},
};
await blocked;
yield {type: 'message_stop'};
} finally {
finalized = true;
}
}

const source = new AnthropicStream(chunks, new AbortController());
const wrapped = meter.trackUsageStreamAnthropic(source as any, 'cus_123');

try {
for await (const _chunk of wrapped) {
break;
}
await new Promise(resolve => setImmediate(resolve));

expect(finalized).toBe(true);
} finally {
release();
}
});

it('should track usage from basic streaming message', async () => {
const meter = createTokenMeter(TEST_API_KEY, config);

Expand Down
94 changes: 94 additions & 0 deletions llm/token-meter/tests/token-meter-openai.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
*/

import Stripe from 'stripe';
import {Stream as OpenAIStream} from 'openai/streaming';
import {createTokenMeter} from '../token-meter';
import type {MeterConfig} from '../types';

Expand Down Expand Up @@ -227,6 +228,99 @@ describe('TokenMeter - OpenAI Provider', () => {
});

describe('Chat Completions - Streaming', () => {
it('does not consume the source before the returned stream is read', async () => {
const meter = createTokenMeter(TEST_API_KEY, config);
let pulls = 0;

async function* chunks() {
pulls += 1;
yield {
id: 'chatcmpl-lazy',
object: 'chat.completion.chunk',
created: Date.now(),
model: 'gpt-4o-mini',
choices: [],
};
}

const source = new OpenAIStream(chunks, new AbortController());
const wrapped = meter.trackUsageStreamOpenAI(source as any, 'cus_123');

await new Promise(resolve => setImmediate(resolve));

expect(pulls).toBe(0);
expect(wrapped).toBeInstanceOf(OpenAIStream);
});

it('closes the source when the returned stream is abandoned', async () => {
const meter = createTokenMeter(TEST_API_KEY, config);
let finalized = false;
let release!: () => void;
const blocked = new Promise<void>(resolve => {
release = resolve;
});

async function* chunks() {
try {
yield {
id: 'chatcmpl-cancel',
object: 'chat.completion.chunk',
created: Date.now(),
model: 'gpt-4o-mini',
choices: [],
};
await blocked;
yield {
id: 'chatcmpl-unused',
object: 'chat.completion.chunk',
created: Date.now(),
model: 'gpt-4o-mini',
choices: [],
};
} finally {
finalized = true;
}
}

const source = new OpenAIStream(chunks, new AbortController());
const wrapped = meter.trackUsageStreamOpenAI(source as any, 'cus_123');

try {
for await (const _chunk of wrapped) {
break;
}
await new Promise(resolve => setImmediate(resolve));

expect(finalized).toBe(true);
} finally {
release();
}
});

it('propagates a source error once and finalizes the source', async () => {
const meter = createTokenMeter(TEST_API_KEY, config);
const sourceError = new Error('source failed before the first chunk');
let finalized = false;

async function* chunks(): AsyncGenerator<any> {
try {
throw sourceError;
} finally {
finalized = true;
}
}

const source = new OpenAIStream(chunks, new AbortController());
const wrapped = meter.trackUsageStreamOpenAI(source as any, 'cus_123');
const iterator = wrapped[Symbol.asyncIterator]();

await expect(iterator.next()).rejects.toBe(sourceError);
await expect(iterator.next()).resolves.toEqual({done: true, value: undefined});
await new Promise(resolve => setImmediate(resolve));

expect(finalized).toBe(true);
});

it('should track usage from basic streaming chat', async () => {
const meter = createTokenMeter(TEST_API_KEY, config);

Expand Down
Loading
Loading