mirror of
https://github.com/Onewon/claude-code.git
synced 2026-04-26 14:51:23 +03:00
Add all the original files from the deleted 0.2.8 npm package
This commit is contained in:
114
vendor/sdk/lib/BetaMessageStream.d.ts
vendored
Normal file
114
vendor/sdk/lib/BetaMessageStream.d.ts
vendored
Normal file
@@ -0,0 +1,114 @@
|
||||
|
||||
import * as Core from '@anthropic-ai/sdk/core';
|
||||
import { AnthropicError, APIUserAbortError } from '@anthropic-ai/sdk/error';
|
||||
import { type BetaContentBlock, Messages as BetaMessages, type BetaMessage, type BetaRawMessageStreamEvent as BetaMessageStreamEvent, type BetaMessageParam, type MessageCreateParams as BetaMessageCreateParams, type MessageCreateParamsBase as BetaMessageCreateParamsBase, type BetaTextCitation } from '@anthropic-ai/sdk/resources/beta/messages/messages';
|
||||
import { type ReadableStream, type Response } from '@anthropic-ai/sdk/_shims/index';
|
||||
export interface MessageStreamEvents {
|
||||
connect: () => void;
|
||||
streamEvent: (event: BetaMessageStreamEvent, snapshot: BetaMessage) => void;
|
||||
text: (textDelta: string, textSnapshot: string) => void;
|
||||
citation: (citation: BetaTextCitation, citationsSnapshot: BetaTextCitation[]) => void;
|
||||
inputJson: (partialJson: string, jsonSnapshot: unknown) => void;
|
||||
thinking: (thinkingDelta: string, thinkingSnapshot: string) => void;
|
||||
message: (message: BetaMessage) => void;
|
||||
contentBlock: (content: BetaContentBlock) => void;
|
||||
finalMessage: (message: BetaMessage) => void;
|
||||
error: (error: AnthropicError) => void;
|
||||
abort: (error: APIUserAbortError) => void;
|
||||
end: () => void;
|
||||
}
|
||||
export declare class BetaMessageStream implements AsyncIterable<BetaMessageStreamEvent> {
|
||||
#private;
|
||||
messages: BetaMessageParam[];
|
||||
receivedMessages: BetaMessage[];
|
||||
controller: AbortController;
|
||||
constructor();
|
||||
get response(): Response | null | undefined;
|
||||
get request_id(): string | null | undefined;
|
||||
/**
|
||||
* Returns the `MessageStream` data, the raw `Response` instance and the ID of the request,
|
||||
* returned vie the `request-id` header which is useful for debugging requests and resporting
|
||||
* issues to Anthropic.
|
||||
*
|
||||
* This is the same as the `APIPromise.withResponse()` method.
|
||||
*
|
||||
* This method will raise an error if you created the stream using `MessageStream.fromReadableStream`
|
||||
* as no `Response` is available.
|
||||
*/
|
||||
withResponse(): Promise<{
|
||||
data: BetaMessageStream;
|
||||
response: Response;
|
||||
request_id: string | null | undefined;
|
||||
}>;
|
||||
/**
|
||||
* Intended for use on the frontend, consuming a stream produced with
|
||||
* `.toReadableStream()` on the backend.
|
||||
*
|
||||
* Note that messages sent to the model do not appear in `.on('message')`
|
||||
* in this context.
|
||||
*/
|
||||
static fromReadableStream(stream: ReadableStream): BetaMessageStream;
|
||||
static createMessage(messages: BetaMessages, params: BetaMessageCreateParamsBase, options?: Core.RequestOptions): BetaMessageStream;
|
||||
protected _run(executor: () => Promise<any>): void;
|
||||
protected _addMessageParam(message: BetaMessageParam): void;
|
||||
protected _addMessage(message: BetaMessage, emit?: boolean): void;
|
||||
protected _createMessage(messages: BetaMessages, params: BetaMessageCreateParams, options?: Core.RequestOptions): Promise<void>;
|
||||
protected _connected(response: Response | null): void;
|
||||
get ended(): boolean;
|
||||
get errored(): boolean;
|
||||
get aborted(): boolean;
|
||||
abort(): void;
|
||||
/**
|
||||
* Adds the listener function to the end of the listeners array for the event.
|
||||
* No checks are made to see if the listener has already been added. Multiple calls passing
|
||||
* the same combination of event and listener will result in the listener being added, and
|
||||
* called, multiple times.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
on<Event extends keyof MessageStreamEvents>(event: Event, listener: MessageStreamEvents[Event]): this;
|
||||
/**
|
||||
* Removes the specified listener from the listener array for the event.
|
||||
* off() will remove, at most, one instance of a listener from the listener array. If any single
|
||||
* listener has been added multiple times to the listener array for the specified event, then
|
||||
* off() must be called multiple times to remove each instance.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
off<Event extends keyof MessageStreamEvents>(event: Event, listener: MessageStreamEvents[Event]): this;
|
||||
/**
|
||||
* Adds a one-time listener function for the event. The next time the event is triggered,
|
||||
* this listener is removed and then invoked.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
once<Event extends keyof MessageStreamEvents>(event: Event, listener: MessageStreamEvents[Event]): this;
|
||||
/**
|
||||
* This is similar to `.once()`, but returns a Promise that resolves the next time
|
||||
* the event is triggered, instead of calling a listener callback.
|
||||
* @returns a Promise that resolves the next time given event is triggered,
|
||||
* or rejects if an error is emitted. (If you request the 'error' event,
|
||||
* returns a promise that resolves with the error).
|
||||
*
|
||||
* Example:
|
||||
*
|
||||
* const message = await stream.emitted('message') // rejects if the stream errors
|
||||
*/
|
||||
emitted<Event extends keyof MessageStreamEvents>(event: Event): Promise<Parameters<MessageStreamEvents[Event]> extends [infer Param] ? Param : Parameters<MessageStreamEvents[Event]> extends [] ? void : Parameters<MessageStreamEvents[Event]>>;
|
||||
done(): Promise<void>;
|
||||
get currentMessage(): BetaMessage | undefined;
|
||||
/**
|
||||
* @returns a promise that resolves with the the final assistant Message response,
|
||||
* or rejects if an error occurred or the stream ended prematurely without producing a Message.
|
||||
*/
|
||||
finalMessage(): Promise<BetaMessage>;
|
||||
/**
|
||||
* @returns a promise that resolves with the the final assistant Message's text response, concatenated
|
||||
* together if there are more than one text blocks.
|
||||
* Rejects if an error occurred or the stream ended prematurely without producing a Message.
|
||||
*/
|
||||
finalText(): Promise<string>;
|
||||
protected _emit<Event extends keyof MessageStreamEvents>(event: Event, ...args: Parameters<MessageStreamEvents[Event]>): void;
|
||||
protected _emitFinal(): void;
|
||||
protected _fromReadableStream(readableStream: ReadableStream, options?: Core.RequestOptions): Promise<void>;
|
||||
[Symbol.asyncIterator](): AsyncIterator<BetaMessageStreamEvent>;
|
||||
toReadableStream(): ReadableStream;
|
||||
}
|
||||
//# sourceMappingURL=BetaMessageStream.d.ts.map
|
||||
1
vendor/sdk/lib/BetaMessageStream.d.ts.map
vendored
Normal file
1
vendor/sdk/lib/BetaMessageStream.d.ts.map
vendored
Normal file
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"BetaMessageStream.d.ts","sourceRoot":"","sources":["../src/lib/BetaMessageStream.ts"],"names":[],"mappings":";AAAA,OAAO,KAAK,IAAI,MAAM,wBAAwB,CAAC;AAC/C,OAAO,EAAE,cAAc,EAAE,iBAAiB,EAAE,MAAM,yBAAyB,CAAC;AAC5E,OAAO,EACL,KAAK,gBAAgB,EACrB,QAAQ,IAAI,YAAY,EACxB,KAAK,WAAW,EAChB,KAAK,yBAAyB,IAAI,sBAAsB,EACxD,KAAK,gBAAgB,EACrB,KAAK,mBAAmB,IAAI,uBAAuB,EACnD,KAAK,uBAAuB,IAAI,2BAA2B,EAE3D,KAAK,gBAAgB,EACtB,MAAM,oDAAoD,CAAC;AAC5D,OAAO,EAAE,KAAK,cAAc,EAAE,KAAK,QAAQ,EAAE,MAAM,gCAAgC,CAAC;AAIpF,MAAM,WAAW,mBAAmB;IAClC,OAAO,EAAE,MAAM,IAAI,CAAC;IACpB,WAAW,EAAE,CAAC,KAAK,EAAE,sBAAsB,EAAE,QAAQ,EAAE,WAAW,KAAK,IAAI,CAAC;IAC5E,IAAI,EAAE,CAAC,SAAS,EAAE,MAAM,EAAE,YAAY,EAAE,MAAM,KAAK,IAAI,CAAC;IACxD,QAAQ,EAAE,CAAC,QAAQ,EAAE,gBAAgB,EAAE,iBAAiB,EAAE,gBAAgB,EAAE,KAAK,IAAI,CAAC;IACtF,SAAS,EAAE,CAAC,WAAW,EAAE,MAAM,EAAE,YAAY,EAAE,OAAO,KAAK,IAAI,CAAC;IAChE,QAAQ,EAAE,CAAC,aAAa,EAAE,MAAM,EAAE,gBAAgB,EAAE,MAAM,KAAK,IAAI,CAAC;IACpE,OAAO,EAAE,CAAC,OAAO,EAAE,WAAW,KAAK,IAAI,CAAC;IACxC,YAAY,EAAE,CAAC,OAAO,EAAE,gBAAgB,KAAK,IAAI,CAAC;IAClD,YAAY,EAAE,CAAC,OAAO,EAAE,WAAW,KAAK,IAAI,CAAC;IAC7C,KAAK,EAAE,CAAC,KAAK,EAAE,cAAc,KAAK,IAAI,CAAC;IACvC,KAAK,EAAE,CAAC,KAAK,EAAE,iBAAiB,KAAK,IAAI,CAAC;IAC1C,GAAG,EAAE,MAAM,IAAI,CAAC;CACjB;AASD,qBAAa,iBAAkB,YAAW,aAAa,CAAC,sBAAsB,CAAC;;IAC7E,QAAQ,EAAE,gBAAgB,EAAE,CAAM;IAClC,gBAAgB,EAAE,WAAW,EAAE,CAAM;IAGrC,UAAU,EAAE,eAAe,CAAyB;;IAsCpD,IAAI,QAAQ,IAAI,QAAQ,GAAG,IAAI,GAAG,SAAS,CAE1C;IAED,IAAI,UAAU,IAAI,MAAM,GAAG,IAAI,GAAG,SAAS,CAE1C;IAED;;;;;;;;;OASG;IACG,YAAY,IAAI,OAAO,CAAC;QAC5B,IAAI,EAAE,iBAAiB,CAAC;QACxB,QAAQ,EAAE,QAAQ,CAAC;QACnB,UAAU,EAAE,MAAM,GAAG,IAAI,GAAG,SAAS,CAAC;KACvC,CAAC;IAaF;;;;;;OAMG;IACH,MAAM,CAAC,kBAAkB,CAAC,MAAM,EAAE,cAAc,GAAG,iBAAiB;IAMpE,MAAM,CAAC,aAAa,CAClB,QAAQ,EAAE,YAAY,EACtB,MAAM,EAAE,2BAA2B,EACnC,OAAO,CAAC,EAAE,IAAI,CAAC,cAAc,GAC5B,iBAAiB;IAepB,SAAS,CAAC,IAAI,CAAC,QAAQ,EAAE,MAAM,OAAO,CAAC,GAAG,CAAC;IAO3C,SAAS,CAAC,gBAAgB,CAAC,OAAO,EAAE,gBAAgB;IAIpD,SAAS,CAAC,WAAW,CAAC,OAAO,EAAE,WAAW,EAAE,IAAI,UAAO;cAOvC,cAAc,CAC5B,QAAQ,EAAE,YAAY,EACtB,MAAM,EAAE,uBAAuB,EAC/B,OAAO,CAAC,EAAE,IAAI,CAAC,cAAc,GAC5B,OAAO,CAAC,IAAI,CAAC;IAoBhB,SAAS,CAAC,UAAU,CAAC,QAAQ,EAAE,QAAQ,GAAG,IAAI;IAQ9C,IAAI,KAAK,IAAI,OAAO,CAEnB;IAED,IAAI,OAAO,IAAI,OAAO,CAErB;IAED,IAAI,OAAO,IAAI,OAAO,CAErB;IAED,KAAK;IAIL;;;;;;OAMG;IACH,EAAE,CAAC,KAAK,SAAS,MAAM,mBAAmB,EAAE,KAAK,EAAE,KAAK,EAAE,QAAQ,EAAE,mBAAmB,CAAC,KAAK,CAAC,GAAG,IAAI;IAOrG;;;;;;OAMG;IACH,GAAG,CAAC,KAAK,SAAS,MAAM,mBAAmB,EAAE,KAAK,EAAE,KAAK,EAAE,QAAQ,EAAE,mBAAmB,CAAC,KAAK,CAAC,GAAG,IAAI;IAQtG;;;;OAIG;IACH,IAAI,CAAC,KAAK,SAAS,MAAM,mBAAmB,EAAE,KAAK,EAAE,KAAK,EAAE,QAAQ,EAAE,mBAAmB,CAAC,KAAK,CAAC,GAAG,IAAI;IAOvG;;;;;;;;;;OAUG;IACH,OAAO,CAAC,KAAK,SAAS,MAAM,mBAAmB,EAC7C,KAAK,EAAE,KAAK,GACX,OAAO,CACR,UAAU,CAAC,mBAAmB,CAAC,KAAK,CAAC,CAAC,SAAS,CAAC,MAAM,KAAK,CAAC,GAAG,KAAK,GAClE,UAAU,CAAC,mBAAmB,CAAC,KAAK,CAAC,CAAC,SAAS,EAAE,GAAG,IAAI,GACxD,UAAU,CAAC,mBAAmB,CAAC,KAAK,CAAC,CAAC,CACzC;IAQK,IAAI,IAAI,OAAO,CAAC,IAAI,CAAC;IAK3B,IAAI,cAAc,IAAI,WAAW,GAAG,SAAS,CAE5C;IASD;;;OAGG;IACG,YAAY,IAAI,OAAO,CAAC,WAAW,CAAC;IAmB1C;;;;OAIG;IACG,SAAS,IAAI,OAAO,CAAC,MAAM,CAAC;IA0BlC,SAAS,CAAC,KAAK,CAAC,KAAK,SAAS,MAAM,mBAAmB,EACrD,KAAK,EAAE,KAAK,EACZ,GAAG,IAAI,EAAE,UAAU,CAAC,mBAAmB,CAAC,KAAK,CAAC,CAAC;IA8CjD,SAAS,CAAC,UAAU;cAmFJ,mBAAmB,CACjC,cAAc,EAAE,cAAc,EAC9B,OAAO,CAAC,EAAE,IAAI,CAAC,cAAc,GAC5B,OAAO,CAAC,IAAI,CAAC;IA2GhB,CAAC,MAAM,CAAC,aAAa,CAAC,IAAI,aAAa,CAAC,sBAAsB,CAAC;IA6D/D,gBAAgB,IAAI,cAAc;CAInC"}
|
||||
547
vendor/sdk/lib/BetaMessageStream.js
vendored
Normal file
547
vendor/sdk/lib/BetaMessageStream.js
vendored
Normal file
@@ -0,0 +1,547 @@
|
||||
"use strict";
|
||||
var __classPrivateFieldSet = (this && this.__classPrivateFieldSet) || function (receiver, state, value, kind, f) {
|
||||
if (kind === "m") throw new TypeError("Private method is not writable");
|
||||
if (kind === "a" && !f) throw new TypeError("Private accessor was defined without a setter");
|
||||
if (typeof state === "function" ? receiver !== state || !f : !state.has(receiver)) throw new TypeError("Cannot write private member to an object whose class did not declare it");
|
||||
return (kind === "a" ? f.call(receiver, value) : f ? f.value = value : state.set(receiver, value)), value;
|
||||
};
|
||||
var __classPrivateFieldGet = (this && this.__classPrivateFieldGet) || function (receiver, state, kind, f) {
|
||||
if (kind === "a" && !f) throw new TypeError("Private accessor was defined without a getter");
|
||||
if (typeof state === "function" ? receiver !== state || !f : !state.has(receiver)) throw new TypeError("Cannot read private member from an object whose class did not declare it");
|
||||
return kind === "m" ? f : kind === "a" ? f.call(receiver) : f ? f.value : state.get(receiver);
|
||||
};
|
||||
var _BetaMessageStream_instances, _BetaMessageStream_currentMessageSnapshot, _BetaMessageStream_connectedPromise, _BetaMessageStream_resolveConnectedPromise, _BetaMessageStream_rejectConnectedPromise, _BetaMessageStream_endPromise, _BetaMessageStream_resolveEndPromise, _BetaMessageStream_rejectEndPromise, _BetaMessageStream_listeners, _BetaMessageStream_ended, _BetaMessageStream_errored, _BetaMessageStream_aborted, _BetaMessageStream_catchingPromiseCreated, _BetaMessageStream_response, _BetaMessageStream_request_id, _BetaMessageStream_getFinalMessage, _BetaMessageStream_getFinalText, _BetaMessageStream_handleError, _BetaMessageStream_beginRequest, _BetaMessageStream_addStreamEvent, _BetaMessageStream_endRequest, _BetaMessageStream_accumulateMessage;
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
exports.BetaMessageStream = void 0;
|
||||
const error_1 = require("@anthropic-ai/sdk/error");
|
||||
const streaming_1 = require("@anthropic-ai/sdk/streaming");
|
||||
const parser_1 = require("../_vendor/partial-json-parser/parser.js");
|
||||
const JSON_BUF_PROPERTY = '__json_buf';
|
||||
class BetaMessageStream {
|
||||
constructor() {
|
||||
_BetaMessageStream_instances.add(this);
|
||||
this.messages = [];
|
||||
this.receivedMessages = [];
|
||||
_BetaMessageStream_currentMessageSnapshot.set(this, void 0);
|
||||
this.controller = new AbortController();
|
||||
_BetaMessageStream_connectedPromise.set(this, void 0);
|
||||
_BetaMessageStream_resolveConnectedPromise.set(this, () => { });
|
||||
_BetaMessageStream_rejectConnectedPromise.set(this, () => { });
|
||||
_BetaMessageStream_endPromise.set(this, void 0);
|
||||
_BetaMessageStream_resolveEndPromise.set(this, () => { });
|
||||
_BetaMessageStream_rejectEndPromise.set(this, () => { });
|
||||
_BetaMessageStream_listeners.set(this, {});
|
||||
_BetaMessageStream_ended.set(this, false);
|
||||
_BetaMessageStream_errored.set(this, false);
|
||||
_BetaMessageStream_aborted.set(this, false);
|
||||
_BetaMessageStream_catchingPromiseCreated.set(this, false);
|
||||
_BetaMessageStream_response.set(this, void 0);
|
||||
_BetaMessageStream_request_id.set(this, void 0);
|
||||
_BetaMessageStream_handleError.set(this, (error) => {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_errored, true, "f");
|
||||
if (error instanceof Error && error.name === 'AbortError') {
|
||||
error = new error_1.APIUserAbortError();
|
||||
}
|
||||
if (error instanceof error_1.APIUserAbortError) {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_aborted, true, "f");
|
||||
return this._emit('abort', error);
|
||||
}
|
||||
if (error instanceof error_1.AnthropicError) {
|
||||
return this._emit('error', error);
|
||||
}
|
||||
if (error instanceof Error) {
|
||||
const anthropicError = new error_1.AnthropicError(error.message);
|
||||
// @ts-ignore
|
||||
anthropicError.cause = error;
|
||||
return this._emit('error', anthropicError);
|
||||
}
|
||||
return this._emit('error', new error_1.AnthropicError(String(error)));
|
||||
});
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_connectedPromise, new Promise((resolve, reject) => {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_resolveConnectedPromise, resolve, "f");
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_rejectConnectedPromise, reject, "f");
|
||||
}), "f");
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_endPromise, new Promise((resolve, reject) => {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_resolveEndPromise, resolve, "f");
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_rejectEndPromise, reject, "f");
|
||||
}), "f");
|
||||
// Don't let these promises cause unhandled rejection errors.
|
||||
// we will manually cause an unhandled rejection error later
|
||||
// if the user hasn't registered any error listener or called
|
||||
// any promise-returning method.
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_connectedPromise, "f").catch(() => { });
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_endPromise, "f").catch(() => { });
|
||||
}
|
||||
get response() {
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_response, "f");
|
||||
}
|
||||
get request_id() {
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_request_id, "f");
|
||||
}
|
||||
/**
|
||||
* Returns the `MessageStream` data, the raw `Response` instance and the ID of the request,
|
||||
* returned vie the `request-id` header which is useful for debugging requests and resporting
|
||||
* issues to Anthropic.
|
||||
*
|
||||
* This is the same as the `APIPromise.withResponse()` method.
|
||||
*
|
||||
* This method will raise an error if you created the stream using `MessageStream.fromReadableStream`
|
||||
* as no `Response` is available.
|
||||
*/
|
||||
async withResponse() {
|
||||
const response = await __classPrivateFieldGet(this, _BetaMessageStream_connectedPromise, "f");
|
||||
if (!response) {
|
||||
throw new Error('Could not resolve a `Response` object');
|
||||
}
|
||||
return {
|
||||
data: this,
|
||||
response,
|
||||
request_id: response.headers.get('request-id'),
|
||||
};
|
||||
}
|
||||
/**
|
||||
* Intended for use on the frontend, consuming a stream produced with
|
||||
* `.toReadableStream()` on the backend.
|
||||
*
|
||||
* Note that messages sent to the model do not appear in `.on('message')`
|
||||
* in this context.
|
||||
*/
|
||||
static fromReadableStream(stream) {
|
||||
const runner = new BetaMessageStream();
|
||||
runner._run(() => runner._fromReadableStream(stream));
|
||||
return runner;
|
||||
}
|
||||
static createMessage(messages, params, options) {
|
||||
const runner = new BetaMessageStream();
|
||||
for (const message of params.messages) {
|
||||
runner._addMessageParam(message);
|
||||
}
|
||||
runner._run(() => runner._createMessage(messages, { ...params, stream: true }, { ...options, headers: { ...options?.headers, 'X-Stainless-Helper-Method': 'stream' } }));
|
||||
return runner;
|
||||
}
|
||||
_run(executor) {
|
||||
executor().then(() => {
|
||||
this._emitFinal();
|
||||
this._emit('end');
|
||||
}, __classPrivateFieldGet(this, _BetaMessageStream_handleError, "f"));
|
||||
}
|
||||
_addMessageParam(message) {
|
||||
this.messages.push(message);
|
||||
}
|
||||
_addMessage(message, emit = true) {
|
||||
this.receivedMessages.push(message);
|
||||
if (emit) {
|
||||
this._emit('message', message);
|
||||
}
|
||||
}
|
||||
async _createMessage(messages, params, options) {
|
||||
const signal = options?.signal;
|
||||
if (signal) {
|
||||
if (signal.aborted)
|
||||
this.controller.abort();
|
||||
signal.addEventListener('abort', () => this.controller.abort());
|
||||
}
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_beginRequest).call(this);
|
||||
const { response, data: stream } = await messages
|
||||
.create({ ...params, stream: true }, { ...options, signal: this.controller.signal })
|
||||
.withResponse();
|
||||
this._connected(response);
|
||||
for await (const event of stream) {
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_addStreamEvent).call(this, event);
|
||||
}
|
||||
if (stream.controller.signal?.aborted) {
|
||||
throw new error_1.APIUserAbortError();
|
||||
}
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_endRequest).call(this);
|
||||
}
|
||||
_connected(response) {
|
||||
if (this.ended)
|
||||
return;
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_response, response, "f");
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_request_id, response?.headers.get('request-id'), "f");
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_resolveConnectedPromise, "f").call(this, response);
|
||||
this._emit('connect');
|
||||
}
|
||||
get ended() {
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_ended, "f");
|
||||
}
|
||||
get errored() {
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_errored, "f");
|
||||
}
|
||||
get aborted() {
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_aborted, "f");
|
||||
}
|
||||
abort() {
|
||||
this.controller.abort();
|
||||
}
|
||||
/**
|
||||
* Adds the listener function to the end of the listeners array for the event.
|
||||
* No checks are made to see if the listener has already been added. Multiple calls passing
|
||||
* the same combination of event and listener will result in the listener being added, and
|
||||
* called, multiple times.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
on(event, listener) {
|
||||
const listeners = __classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event] || (__classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event] = []);
|
||||
listeners.push({ listener });
|
||||
return this;
|
||||
}
|
||||
/**
|
||||
* Removes the specified listener from the listener array for the event.
|
||||
* off() will remove, at most, one instance of a listener from the listener array. If any single
|
||||
* listener has been added multiple times to the listener array for the specified event, then
|
||||
* off() must be called multiple times to remove each instance.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
off(event, listener) {
|
||||
const listeners = __classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event];
|
||||
if (!listeners)
|
||||
return this;
|
||||
const index = listeners.findIndex((l) => l.listener === listener);
|
||||
if (index >= 0)
|
||||
listeners.splice(index, 1);
|
||||
return this;
|
||||
}
|
||||
/**
|
||||
* Adds a one-time listener function for the event. The next time the event is triggered,
|
||||
* this listener is removed and then invoked.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
once(event, listener) {
|
||||
const listeners = __classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event] || (__classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event] = []);
|
||||
listeners.push({ listener, once: true });
|
||||
return this;
|
||||
}
|
||||
/**
|
||||
* This is similar to `.once()`, but returns a Promise that resolves the next time
|
||||
* the event is triggered, instead of calling a listener callback.
|
||||
* @returns a Promise that resolves the next time given event is triggered,
|
||||
* or rejects if an error is emitted. (If you request the 'error' event,
|
||||
* returns a promise that resolves with the error).
|
||||
*
|
||||
* Example:
|
||||
*
|
||||
* const message = await stream.emitted('message') // rejects if the stream errors
|
||||
*/
|
||||
emitted(event) {
|
||||
return new Promise((resolve, reject) => {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_catchingPromiseCreated, true, "f");
|
||||
if (event !== 'error')
|
||||
this.once('error', reject);
|
||||
this.once(event, resolve);
|
||||
});
|
||||
}
|
||||
async done() {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_catchingPromiseCreated, true, "f");
|
||||
await __classPrivateFieldGet(this, _BetaMessageStream_endPromise, "f");
|
||||
}
|
||||
get currentMessage() {
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_currentMessageSnapshot, "f");
|
||||
}
|
||||
/**
|
||||
* @returns a promise that resolves with the the final assistant Message response,
|
||||
* or rejects if an error occurred or the stream ended prematurely without producing a Message.
|
||||
*/
|
||||
async finalMessage() {
|
||||
await this.done();
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_getFinalMessage).call(this);
|
||||
}
|
||||
/**
|
||||
* @returns a promise that resolves with the the final assistant Message's text response, concatenated
|
||||
* together if there are more than one text blocks.
|
||||
* Rejects if an error occurred or the stream ended prematurely without producing a Message.
|
||||
*/
|
||||
async finalText() {
|
||||
await this.done();
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_getFinalText).call(this);
|
||||
}
|
||||
_emit(event, ...args) {
|
||||
// make sure we don't emit any MessageStreamEvents after end
|
||||
if (__classPrivateFieldGet(this, _BetaMessageStream_ended, "f"))
|
||||
return;
|
||||
if (event === 'end') {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_ended, true, "f");
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_resolveEndPromise, "f").call(this);
|
||||
}
|
||||
const listeners = __classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event];
|
||||
if (listeners) {
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event] = listeners.filter((l) => !l.once);
|
||||
listeners.forEach(({ listener }) => listener(...args));
|
||||
}
|
||||
if (event === 'abort') {
|
||||
const error = args[0];
|
||||
if (!__classPrivateFieldGet(this, _BetaMessageStream_catchingPromiseCreated, "f") && !listeners?.length) {
|
||||
Promise.reject(error);
|
||||
}
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_rejectConnectedPromise, "f").call(this, error);
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_rejectEndPromise, "f").call(this, error);
|
||||
this._emit('end');
|
||||
return;
|
||||
}
|
||||
if (event === 'error') {
|
||||
// NOTE: _emit('error', error) should only be called from #handleError().
|
||||
const error = args[0];
|
||||
if (!__classPrivateFieldGet(this, _BetaMessageStream_catchingPromiseCreated, "f") && !listeners?.length) {
|
||||
// Trigger an unhandled rejection if the user hasn't registered any error handlers.
|
||||
// If you are seeing stack traces here, make sure to handle errors via either:
|
||||
// - runner.on('error', () => ...)
|
||||
// - await runner.done()
|
||||
// - await runner.final...()
|
||||
// - etc.
|
||||
Promise.reject(error);
|
||||
}
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_rejectConnectedPromise, "f").call(this, error);
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_rejectEndPromise, "f").call(this, error);
|
||||
this._emit('end');
|
||||
}
|
||||
}
|
||||
_emitFinal() {
|
||||
const finalMessage = this.receivedMessages.at(-1);
|
||||
if (finalMessage) {
|
||||
this._emit('finalMessage', __classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_getFinalMessage).call(this));
|
||||
}
|
||||
}
|
||||
async _fromReadableStream(readableStream, options) {
|
||||
const signal = options?.signal;
|
||||
if (signal) {
|
||||
if (signal.aborted)
|
||||
this.controller.abort();
|
||||
signal.addEventListener('abort', () => this.controller.abort());
|
||||
}
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_beginRequest).call(this);
|
||||
this._connected(null);
|
||||
const stream = streaming_1.Stream.fromReadableStream(readableStream, this.controller);
|
||||
for await (const event of stream) {
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_addStreamEvent).call(this, event);
|
||||
}
|
||||
if (stream.controller.signal?.aborted) {
|
||||
throw new error_1.APIUserAbortError();
|
||||
}
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_endRequest).call(this);
|
||||
}
|
||||
[(_BetaMessageStream_currentMessageSnapshot = new WeakMap(), _BetaMessageStream_connectedPromise = new WeakMap(), _BetaMessageStream_resolveConnectedPromise = new WeakMap(), _BetaMessageStream_rejectConnectedPromise = new WeakMap(), _BetaMessageStream_endPromise = new WeakMap(), _BetaMessageStream_resolveEndPromise = new WeakMap(), _BetaMessageStream_rejectEndPromise = new WeakMap(), _BetaMessageStream_listeners = new WeakMap(), _BetaMessageStream_ended = new WeakMap(), _BetaMessageStream_errored = new WeakMap(), _BetaMessageStream_aborted = new WeakMap(), _BetaMessageStream_catchingPromiseCreated = new WeakMap(), _BetaMessageStream_response = new WeakMap(), _BetaMessageStream_request_id = new WeakMap(), _BetaMessageStream_handleError = new WeakMap(), _BetaMessageStream_instances = new WeakSet(), _BetaMessageStream_getFinalMessage = function _BetaMessageStream_getFinalMessage() {
|
||||
if (this.receivedMessages.length === 0) {
|
||||
throw new error_1.AnthropicError('stream ended without producing a Message with role=assistant');
|
||||
}
|
||||
return this.receivedMessages.at(-1);
|
||||
}, _BetaMessageStream_getFinalText = function _BetaMessageStream_getFinalText() {
|
||||
if (this.receivedMessages.length === 0) {
|
||||
throw new error_1.AnthropicError('stream ended without producing a Message with role=assistant');
|
||||
}
|
||||
const textBlocks = this.receivedMessages
|
||||
.at(-1)
|
||||
.content.filter((block) => block.type === 'text')
|
||||
.map((block) => block.text);
|
||||
if (textBlocks.length === 0) {
|
||||
throw new error_1.AnthropicError('stream ended without producing a content block with type=text');
|
||||
}
|
||||
return textBlocks.join(' ');
|
||||
}, _BetaMessageStream_beginRequest = function _BetaMessageStream_beginRequest() {
|
||||
if (this.ended)
|
||||
return;
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_currentMessageSnapshot, undefined, "f");
|
||||
}, _BetaMessageStream_addStreamEvent = function _BetaMessageStream_addStreamEvent(event) {
|
||||
if (this.ended)
|
||||
return;
|
||||
const messageSnapshot = __classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_accumulateMessage).call(this, event);
|
||||
this._emit('streamEvent', event, messageSnapshot);
|
||||
switch (event.type) {
|
||||
case 'content_block_delta': {
|
||||
const content = messageSnapshot.content.at(-1);
|
||||
switch (event.delta.type) {
|
||||
case 'text_delta': {
|
||||
if (content.type === 'text') {
|
||||
this._emit('text', event.delta.text, content.text || '');
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'citations_delta': {
|
||||
if (content.type === 'text') {
|
||||
this._emit('citation', event.delta.citation, content.citations ?? []);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'input_json_delta': {
|
||||
if (content.type === 'tool_use' && content.input) {
|
||||
this._emit('inputJson', event.delta.partial_json, content.input);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'thinking_delta': {
|
||||
if (content.type === 'thinking') {
|
||||
this._emit('thinking', event.delta.thinking, content.thinking);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'signature_delta': {
|
||||
// we don't emit anything special in this case.
|
||||
break;
|
||||
}
|
||||
default:
|
||||
checkNever(event.delta);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'message_stop': {
|
||||
this._addMessageParam(messageSnapshot);
|
||||
this._addMessage(messageSnapshot, true);
|
||||
break;
|
||||
}
|
||||
case 'content_block_stop': {
|
||||
this._emit('contentBlock', messageSnapshot.content.at(-1));
|
||||
break;
|
||||
}
|
||||
case 'message_start': {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_currentMessageSnapshot, messageSnapshot, "f");
|
||||
break;
|
||||
}
|
||||
case 'content_block_start':
|
||||
case 'message_delta':
|
||||
break;
|
||||
}
|
||||
}, _BetaMessageStream_endRequest = function _BetaMessageStream_endRequest() {
|
||||
if (this.ended) {
|
||||
throw new error_1.AnthropicError(`stream has ended, this shouldn't happen`);
|
||||
}
|
||||
const snapshot = __classPrivateFieldGet(this, _BetaMessageStream_currentMessageSnapshot, "f");
|
||||
if (!snapshot) {
|
||||
throw new error_1.AnthropicError(`request ended without sending any chunks`);
|
||||
}
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_currentMessageSnapshot, undefined, "f");
|
||||
return snapshot;
|
||||
}, _BetaMessageStream_accumulateMessage = function _BetaMessageStream_accumulateMessage(event) {
|
||||
let snapshot = __classPrivateFieldGet(this, _BetaMessageStream_currentMessageSnapshot, "f");
|
||||
if (event.type === 'message_start') {
|
||||
if (snapshot) {
|
||||
throw new error_1.AnthropicError(`Unexpected event order, got ${event.type} before receiving "message_stop"`);
|
||||
}
|
||||
return event.message;
|
||||
}
|
||||
if (!snapshot) {
|
||||
throw new error_1.AnthropicError(`Unexpected event order, got ${event.type} before "message_start"`);
|
||||
}
|
||||
switch (event.type) {
|
||||
case 'message_stop':
|
||||
return snapshot;
|
||||
case 'message_delta':
|
||||
snapshot.stop_reason = event.delta.stop_reason;
|
||||
snapshot.stop_sequence = event.delta.stop_sequence;
|
||||
snapshot.usage.output_tokens = event.usage.output_tokens;
|
||||
return snapshot;
|
||||
case 'content_block_start':
|
||||
snapshot.content.push(event.content_block);
|
||||
return snapshot;
|
||||
case 'content_block_delta': {
|
||||
const snapshotContent = snapshot.content.at(event.index);
|
||||
switch (event.delta.type) {
|
||||
case 'text_delta': {
|
||||
if (snapshotContent?.type === 'text') {
|
||||
snapshotContent.text += event.delta.text;
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'citations_delta': {
|
||||
if (snapshotContent?.type === 'text') {
|
||||
snapshotContent.citations ?? (snapshotContent.citations = []);
|
||||
snapshotContent.citations.push(event.delta.citation);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'input_json_delta': {
|
||||
if (snapshotContent?.type === 'tool_use') {
|
||||
// we need to keep track of the raw JSON string as well so that we can
|
||||
// re-parse it for each delta, for now we just store it as an untyped
|
||||
// non-enumerable property on the snapshot
|
||||
let jsonBuf = snapshotContent[JSON_BUF_PROPERTY] || '';
|
||||
jsonBuf += event.delta.partial_json;
|
||||
Object.defineProperty(snapshotContent, JSON_BUF_PROPERTY, {
|
||||
value: jsonBuf,
|
||||
enumerable: false,
|
||||
writable: true,
|
||||
});
|
||||
if (jsonBuf) {
|
||||
snapshotContent.input = (0, parser_1.partialParse)(jsonBuf);
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'thinking_delta': {
|
||||
if (snapshotContent?.type === 'thinking') {
|
||||
snapshotContent.thinking += event.delta.thinking;
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'signature_delta': {
|
||||
if (snapshotContent?.type === 'thinking') {
|
||||
snapshotContent.signature += event.delta.signature;
|
||||
}
|
||||
break;
|
||||
}
|
||||
default:
|
||||
checkNever(event.delta);
|
||||
}
|
||||
return snapshot;
|
||||
}
|
||||
case 'content_block_stop':
|
||||
return snapshot;
|
||||
}
|
||||
}, Symbol.asyncIterator)]() {
|
||||
const pushQueue = [];
|
||||
const readQueue = [];
|
||||
let done = false;
|
||||
this.on('streamEvent', (event) => {
|
||||
const reader = readQueue.shift();
|
||||
if (reader) {
|
||||
reader.resolve(event);
|
||||
}
|
||||
else {
|
||||
pushQueue.push(event);
|
||||
}
|
||||
});
|
||||
this.on('end', () => {
|
||||
done = true;
|
||||
for (const reader of readQueue) {
|
||||
reader.resolve(undefined);
|
||||
}
|
||||
readQueue.length = 0;
|
||||
});
|
||||
this.on('abort', (err) => {
|
||||
done = true;
|
||||
for (const reader of readQueue) {
|
||||
reader.reject(err);
|
||||
}
|
||||
readQueue.length = 0;
|
||||
});
|
||||
this.on('error', (err) => {
|
||||
done = true;
|
||||
for (const reader of readQueue) {
|
||||
reader.reject(err);
|
||||
}
|
||||
readQueue.length = 0;
|
||||
});
|
||||
return {
|
||||
next: async () => {
|
||||
if (!pushQueue.length) {
|
||||
if (done) {
|
||||
return { value: undefined, done: true };
|
||||
}
|
||||
return new Promise((resolve, reject) => readQueue.push({ resolve, reject })).then((chunk) => (chunk ? { value: chunk, done: false } : { value: undefined, done: true }));
|
||||
}
|
||||
const chunk = pushQueue.shift();
|
||||
return { value: chunk, done: false };
|
||||
},
|
||||
return: async () => {
|
||||
this.abort();
|
||||
return { value: undefined, done: true };
|
||||
},
|
||||
};
|
||||
}
|
||||
toReadableStream() {
|
||||
const stream = new streaming_1.Stream(this[Symbol.asyncIterator].bind(this), this.controller);
|
||||
return stream.toReadableStream();
|
||||
}
|
||||
}
|
||||
exports.BetaMessageStream = BetaMessageStream;
|
||||
// used to ensure exhaustive case matching without throwing a runtime error
|
||||
function checkNever(x) { }
|
||||
//# sourceMappingURL=BetaMessageStream.js.map
|
||||
1
vendor/sdk/lib/BetaMessageStream.js.map
vendored
Normal file
1
vendor/sdk/lib/BetaMessageStream.js.map
vendored
Normal file
File diff suppressed because one or more lines are too long
543
vendor/sdk/lib/BetaMessageStream.mjs
vendored
Normal file
543
vendor/sdk/lib/BetaMessageStream.mjs
vendored
Normal file
@@ -0,0 +1,543 @@
|
||||
var __classPrivateFieldSet = (this && this.__classPrivateFieldSet) || function (receiver, state, value, kind, f) {
|
||||
if (kind === "m") throw new TypeError("Private method is not writable");
|
||||
if (kind === "a" && !f) throw new TypeError("Private accessor was defined without a setter");
|
||||
if (typeof state === "function" ? receiver !== state || !f : !state.has(receiver)) throw new TypeError("Cannot write private member to an object whose class did not declare it");
|
||||
return (kind === "a" ? f.call(receiver, value) : f ? f.value = value : state.set(receiver, value)), value;
|
||||
};
|
||||
var __classPrivateFieldGet = (this && this.__classPrivateFieldGet) || function (receiver, state, kind, f) {
|
||||
if (kind === "a" && !f) throw new TypeError("Private accessor was defined without a getter");
|
||||
if (typeof state === "function" ? receiver !== state || !f : !state.has(receiver)) throw new TypeError("Cannot read private member from an object whose class did not declare it");
|
||||
return kind === "m" ? f : kind === "a" ? f.call(receiver) : f ? f.value : state.get(receiver);
|
||||
};
|
||||
var _BetaMessageStream_instances, _BetaMessageStream_currentMessageSnapshot, _BetaMessageStream_connectedPromise, _BetaMessageStream_resolveConnectedPromise, _BetaMessageStream_rejectConnectedPromise, _BetaMessageStream_endPromise, _BetaMessageStream_resolveEndPromise, _BetaMessageStream_rejectEndPromise, _BetaMessageStream_listeners, _BetaMessageStream_ended, _BetaMessageStream_errored, _BetaMessageStream_aborted, _BetaMessageStream_catchingPromiseCreated, _BetaMessageStream_response, _BetaMessageStream_request_id, _BetaMessageStream_getFinalMessage, _BetaMessageStream_getFinalText, _BetaMessageStream_handleError, _BetaMessageStream_beginRequest, _BetaMessageStream_addStreamEvent, _BetaMessageStream_endRequest, _BetaMessageStream_accumulateMessage;
|
||||
import { AnthropicError, APIUserAbortError } from '@anthropic-ai/sdk/error';
|
||||
import { Stream } from '@anthropic-ai/sdk/streaming';
|
||||
import { partialParse } from "../_vendor/partial-json-parser/parser.mjs";
|
||||
const JSON_BUF_PROPERTY = '__json_buf';
|
||||
export class BetaMessageStream {
|
||||
constructor() {
|
||||
_BetaMessageStream_instances.add(this);
|
||||
this.messages = [];
|
||||
this.receivedMessages = [];
|
||||
_BetaMessageStream_currentMessageSnapshot.set(this, void 0);
|
||||
this.controller = new AbortController();
|
||||
_BetaMessageStream_connectedPromise.set(this, void 0);
|
||||
_BetaMessageStream_resolveConnectedPromise.set(this, () => { });
|
||||
_BetaMessageStream_rejectConnectedPromise.set(this, () => { });
|
||||
_BetaMessageStream_endPromise.set(this, void 0);
|
||||
_BetaMessageStream_resolveEndPromise.set(this, () => { });
|
||||
_BetaMessageStream_rejectEndPromise.set(this, () => { });
|
||||
_BetaMessageStream_listeners.set(this, {});
|
||||
_BetaMessageStream_ended.set(this, false);
|
||||
_BetaMessageStream_errored.set(this, false);
|
||||
_BetaMessageStream_aborted.set(this, false);
|
||||
_BetaMessageStream_catchingPromiseCreated.set(this, false);
|
||||
_BetaMessageStream_response.set(this, void 0);
|
||||
_BetaMessageStream_request_id.set(this, void 0);
|
||||
_BetaMessageStream_handleError.set(this, (error) => {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_errored, true, "f");
|
||||
if (error instanceof Error && error.name === 'AbortError') {
|
||||
error = new APIUserAbortError();
|
||||
}
|
||||
if (error instanceof APIUserAbortError) {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_aborted, true, "f");
|
||||
return this._emit('abort', error);
|
||||
}
|
||||
if (error instanceof AnthropicError) {
|
||||
return this._emit('error', error);
|
||||
}
|
||||
if (error instanceof Error) {
|
||||
const anthropicError = new AnthropicError(error.message);
|
||||
// @ts-ignore
|
||||
anthropicError.cause = error;
|
||||
return this._emit('error', anthropicError);
|
||||
}
|
||||
return this._emit('error', new AnthropicError(String(error)));
|
||||
});
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_connectedPromise, new Promise((resolve, reject) => {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_resolveConnectedPromise, resolve, "f");
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_rejectConnectedPromise, reject, "f");
|
||||
}), "f");
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_endPromise, new Promise((resolve, reject) => {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_resolveEndPromise, resolve, "f");
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_rejectEndPromise, reject, "f");
|
||||
}), "f");
|
||||
// Don't let these promises cause unhandled rejection errors.
|
||||
// we will manually cause an unhandled rejection error later
|
||||
// if the user hasn't registered any error listener or called
|
||||
// any promise-returning method.
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_connectedPromise, "f").catch(() => { });
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_endPromise, "f").catch(() => { });
|
||||
}
|
||||
get response() {
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_response, "f");
|
||||
}
|
||||
get request_id() {
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_request_id, "f");
|
||||
}
|
||||
/**
|
||||
* Returns the `MessageStream` data, the raw `Response` instance and the ID of the request,
|
||||
* returned vie the `request-id` header which is useful for debugging requests and resporting
|
||||
* issues to Anthropic.
|
||||
*
|
||||
* This is the same as the `APIPromise.withResponse()` method.
|
||||
*
|
||||
* This method will raise an error if you created the stream using `MessageStream.fromReadableStream`
|
||||
* as no `Response` is available.
|
||||
*/
|
||||
async withResponse() {
|
||||
const response = await __classPrivateFieldGet(this, _BetaMessageStream_connectedPromise, "f");
|
||||
if (!response) {
|
||||
throw new Error('Could not resolve a `Response` object');
|
||||
}
|
||||
return {
|
||||
data: this,
|
||||
response,
|
||||
request_id: response.headers.get('request-id'),
|
||||
};
|
||||
}
|
||||
/**
|
||||
* Intended for use on the frontend, consuming a stream produced with
|
||||
* `.toReadableStream()` on the backend.
|
||||
*
|
||||
* Note that messages sent to the model do not appear in `.on('message')`
|
||||
* in this context.
|
||||
*/
|
||||
static fromReadableStream(stream) {
|
||||
const runner = new BetaMessageStream();
|
||||
runner._run(() => runner._fromReadableStream(stream));
|
||||
return runner;
|
||||
}
|
||||
static createMessage(messages, params, options) {
|
||||
const runner = new BetaMessageStream();
|
||||
for (const message of params.messages) {
|
||||
runner._addMessageParam(message);
|
||||
}
|
||||
runner._run(() => runner._createMessage(messages, { ...params, stream: true }, { ...options, headers: { ...options?.headers, 'X-Stainless-Helper-Method': 'stream' } }));
|
||||
return runner;
|
||||
}
|
||||
_run(executor) {
|
||||
executor().then(() => {
|
||||
this._emitFinal();
|
||||
this._emit('end');
|
||||
}, __classPrivateFieldGet(this, _BetaMessageStream_handleError, "f"));
|
||||
}
|
||||
_addMessageParam(message) {
|
||||
this.messages.push(message);
|
||||
}
|
||||
_addMessage(message, emit = true) {
|
||||
this.receivedMessages.push(message);
|
||||
if (emit) {
|
||||
this._emit('message', message);
|
||||
}
|
||||
}
|
||||
async _createMessage(messages, params, options) {
|
||||
const signal = options?.signal;
|
||||
if (signal) {
|
||||
if (signal.aborted)
|
||||
this.controller.abort();
|
||||
signal.addEventListener('abort', () => this.controller.abort());
|
||||
}
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_beginRequest).call(this);
|
||||
const { response, data: stream } = await messages
|
||||
.create({ ...params, stream: true }, { ...options, signal: this.controller.signal })
|
||||
.withResponse();
|
||||
this._connected(response);
|
||||
for await (const event of stream) {
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_addStreamEvent).call(this, event);
|
||||
}
|
||||
if (stream.controller.signal?.aborted) {
|
||||
throw new APIUserAbortError();
|
||||
}
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_endRequest).call(this);
|
||||
}
|
||||
_connected(response) {
|
||||
if (this.ended)
|
||||
return;
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_response, response, "f");
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_request_id, response?.headers.get('request-id'), "f");
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_resolveConnectedPromise, "f").call(this, response);
|
||||
this._emit('connect');
|
||||
}
|
||||
get ended() {
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_ended, "f");
|
||||
}
|
||||
get errored() {
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_errored, "f");
|
||||
}
|
||||
get aborted() {
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_aborted, "f");
|
||||
}
|
||||
abort() {
|
||||
this.controller.abort();
|
||||
}
|
||||
/**
|
||||
* Adds the listener function to the end of the listeners array for the event.
|
||||
* No checks are made to see if the listener has already been added. Multiple calls passing
|
||||
* the same combination of event and listener will result in the listener being added, and
|
||||
* called, multiple times.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
on(event, listener) {
|
||||
const listeners = __classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event] || (__classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event] = []);
|
||||
listeners.push({ listener });
|
||||
return this;
|
||||
}
|
||||
/**
|
||||
* Removes the specified listener from the listener array for the event.
|
||||
* off() will remove, at most, one instance of a listener from the listener array. If any single
|
||||
* listener has been added multiple times to the listener array for the specified event, then
|
||||
* off() must be called multiple times to remove each instance.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
off(event, listener) {
|
||||
const listeners = __classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event];
|
||||
if (!listeners)
|
||||
return this;
|
||||
const index = listeners.findIndex((l) => l.listener === listener);
|
||||
if (index >= 0)
|
||||
listeners.splice(index, 1);
|
||||
return this;
|
||||
}
|
||||
/**
|
||||
* Adds a one-time listener function for the event. The next time the event is triggered,
|
||||
* this listener is removed and then invoked.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
once(event, listener) {
|
||||
const listeners = __classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event] || (__classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event] = []);
|
||||
listeners.push({ listener, once: true });
|
||||
return this;
|
||||
}
|
||||
/**
|
||||
* This is similar to `.once()`, but returns a Promise that resolves the next time
|
||||
* the event is triggered, instead of calling a listener callback.
|
||||
* @returns a Promise that resolves the next time given event is triggered,
|
||||
* or rejects if an error is emitted. (If you request the 'error' event,
|
||||
* returns a promise that resolves with the error).
|
||||
*
|
||||
* Example:
|
||||
*
|
||||
* const message = await stream.emitted('message') // rejects if the stream errors
|
||||
*/
|
||||
emitted(event) {
|
||||
return new Promise((resolve, reject) => {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_catchingPromiseCreated, true, "f");
|
||||
if (event !== 'error')
|
||||
this.once('error', reject);
|
||||
this.once(event, resolve);
|
||||
});
|
||||
}
|
||||
async done() {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_catchingPromiseCreated, true, "f");
|
||||
await __classPrivateFieldGet(this, _BetaMessageStream_endPromise, "f");
|
||||
}
|
||||
get currentMessage() {
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_currentMessageSnapshot, "f");
|
||||
}
|
||||
/**
|
||||
* @returns a promise that resolves with the the final assistant Message response,
|
||||
* or rejects if an error occurred or the stream ended prematurely without producing a Message.
|
||||
*/
|
||||
async finalMessage() {
|
||||
await this.done();
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_getFinalMessage).call(this);
|
||||
}
|
||||
/**
|
||||
* @returns a promise that resolves with the the final assistant Message's text response, concatenated
|
||||
* together if there are more than one text blocks.
|
||||
* Rejects if an error occurred or the stream ended prematurely without producing a Message.
|
||||
*/
|
||||
async finalText() {
|
||||
await this.done();
|
||||
return __classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_getFinalText).call(this);
|
||||
}
|
||||
_emit(event, ...args) {
|
||||
// make sure we don't emit any MessageStreamEvents after end
|
||||
if (__classPrivateFieldGet(this, _BetaMessageStream_ended, "f"))
|
||||
return;
|
||||
if (event === 'end') {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_ended, true, "f");
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_resolveEndPromise, "f").call(this);
|
||||
}
|
||||
const listeners = __classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event];
|
||||
if (listeners) {
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_listeners, "f")[event] = listeners.filter((l) => !l.once);
|
||||
listeners.forEach(({ listener }) => listener(...args));
|
||||
}
|
||||
if (event === 'abort') {
|
||||
const error = args[0];
|
||||
if (!__classPrivateFieldGet(this, _BetaMessageStream_catchingPromiseCreated, "f") && !listeners?.length) {
|
||||
Promise.reject(error);
|
||||
}
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_rejectConnectedPromise, "f").call(this, error);
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_rejectEndPromise, "f").call(this, error);
|
||||
this._emit('end');
|
||||
return;
|
||||
}
|
||||
if (event === 'error') {
|
||||
// NOTE: _emit('error', error) should only be called from #handleError().
|
||||
const error = args[0];
|
||||
if (!__classPrivateFieldGet(this, _BetaMessageStream_catchingPromiseCreated, "f") && !listeners?.length) {
|
||||
// Trigger an unhandled rejection if the user hasn't registered any error handlers.
|
||||
// If you are seeing stack traces here, make sure to handle errors via either:
|
||||
// - runner.on('error', () => ...)
|
||||
// - await runner.done()
|
||||
// - await runner.final...()
|
||||
// - etc.
|
||||
Promise.reject(error);
|
||||
}
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_rejectConnectedPromise, "f").call(this, error);
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_rejectEndPromise, "f").call(this, error);
|
||||
this._emit('end');
|
||||
}
|
||||
}
|
||||
_emitFinal() {
|
||||
const finalMessage = this.receivedMessages.at(-1);
|
||||
if (finalMessage) {
|
||||
this._emit('finalMessage', __classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_getFinalMessage).call(this));
|
||||
}
|
||||
}
|
||||
async _fromReadableStream(readableStream, options) {
|
||||
const signal = options?.signal;
|
||||
if (signal) {
|
||||
if (signal.aborted)
|
||||
this.controller.abort();
|
||||
signal.addEventListener('abort', () => this.controller.abort());
|
||||
}
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_beginRequest).call(this);
|
||||
this._connected(null);
|
||||
const stream = Stream.fromReadableStream(readableStream, this.controller);
|
||||
for await (const event of stream) {
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_addStreamEvent).call(this, event);
|
||||
}
|
||||
if (stream.controller.signal?.aborted) {
|
||||
throw new APIUserAbortError();
|
||||
}
|
||||
__classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_endRequest).call(this);
|
||||
}
|
||||
[(_BetaMessageStream_currentMessageSnapshot = new WeakMap(), _BetaMessageStream_connectedPromise = new WeakMap(), _BetaMessageStream_resolveConnectedPromise = new WeakMap(), _BetaMessageStream_rejectConnectedPromise = new WeakMap(), _BetaMessageStream_endPromise = new WeakMap(), _BetaMessageStream_resolveEndPromise = new WeakMap(), _BetaMessageStream_rejectEndPromise = new WeakMap(), _BetaMessageStream_listeners = new WeakMap(), _BetaMessageStream_ended = new WeakMap(), _BetaMessageStream_errored = new WeakMap(), _BetaMessageStream_aborted = new WeakMap(), _BetaMessageStream_catchingPromiseCreated = new WeakMap(), _BetaMessageStream_response = new WeakMap(), _BetaMessageStream_request_id = new WeakMap(), _BetaMessageStream_handleError = new WeakMap(), _BetaMessageStream_instances = new WeakSet(), _BetaMessageStream_getFinalMessage = function _BetaMessageStream_getFinalMessage() {
|
||||
if (this.receivedMessages.length === 0) {
|
||||
throw new AnthropicError('stream ended without producing a Message with role=assistant');
|
||||
}
|
||||
return this.receivedMessages.at(-1);
|
||||
}, _BetaMessageStream_getFinalText = function _BetaMessageStream_getFinalText() {
|
||||
if (this.receivedMessages.length === 0) {
|
||||
throw new AnthropicError('stream ended without producing a Message with role=assistant');
|
||||
}
|
||||
const textBlocks = this.receivedMessages
|
||||
.at(-1)
|
||||
.content.filter((block) => block.type === 'text')
|
||||
.map((block) => block.text);
|
||||
if (textBlocks.length === 0) {
|
||||
throw new AnthropicError('stream ended without producing a content block with type=text');
|
||||
}
|
||||
return textBlocks.join(' ');
|
||||
}, _BetaMessageStream_beginRequest = function _BetaMessageStream_beginRequest() {
|
||||
if (this.ended)
|
||||
return;
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_currentMessageSnapshot, undefined, "f");
|
||||
}, _BetaMessageStream_addStreamEvent = function _BetaMessageStream_addStreamEvent(event) {
|
||||
if (this.ended)
|
||||
return;
|
||||
const messageSnapshot = __classPrivateFieldGet(this, _BetaMessageStream_instances, "m", _BetaMessageStream_accumulateMessage).call(this, event);
|
||||
this._emit('streamEvent', event, messageSnapshot);
|
||||
switch (event.type) {
|
||||
case 'content_block_delta': {
|
||||
const content = messageSnapshot.content.at(-1);
|
||||
switch (event.delta.type) {
|
||||
case 'text_delta': {
|
||||
if (content.type === 'text') {
|
||||
this._emit('text', event.delta.text, content.text || '');
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'citations_delta': {
|
||||
if (content.type === 'text') {
|
||||
this._emit('citation', event.delta.citation, content.citations ?? []);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'input_json_delta': {
|
||||
if (content.type === 'tool_use' && content.input) {
|
||||
this._emit('inputJson', event.delta.partial_json, content.input);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'thinking_delta': {
|
||||
if (content.type === 'thinking') {
|
||||
this._emit('thinking', event.delta.thinking, content.thinking);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'signature_delta': {
|
||||
// we don't emit anything special in this case.
|
||||
break;
|
||||
}
|
||||
default:
|
||||
checkNever(event.delta);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'message_stop': {
|
||||
this._addMessageParam(messageSnapshot);
|
||||
this._addMessage(messageSnapshot, true);
|
||||
break;
|
||||
}
|
||||
case 'content_block_stop': {
|
||||
this._emit('contentBlock', messageSnapshot.content.at(-1));
|
||||
break;
|
||||
}
|
||||
case 'message_start': {
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_currentMessageSnapshot, messageSnapshot, "f");
|
||||
break;
|
||||
}
|
||||
case 'content_block_start':
|
||||
case 'message_delta':
|
||||
break;
|
||||
}
|
||||
}, _BetaMessageStream_endRequest = function _BetaMessageStream_endRequest() {
|
||||
if (this.ended) {
|
||||
throw new AnthropicError(`stream has ended, this shouldn't happen`);
|
||||
}
|
||||
const snapshot = __classPrivateFieldGet(this, _BetaMessageStream_currentMessageSnapshot, "f");
|
||||
if (!snapshot) {
|
||||
throw new AnthropicError(`request ended without sending any chunks`);
|
||||
}
|
||||
__classPrivateFieldSet(this, _BetaMessageStream_currentMessageSnapshot, undefined, "f");
|
||||
return snapshot;
|
||||
}, _BetaMessageStream_accumulateMessage = function _BetaMessageStream_accumulateMessage(event) {
|
||||
let snapshot = __classPrivateFieldGet(this, _BetaMessageStream_currentMessageSnapshot, "f");
|
||||
if (event.type === 'message_start') {
|
||||
if (snapshot) {
|
||||
throw new AnthropicError(`Unexpected event order, got ${event.type} before receiving "message_stop"`);
|
||||
}
|
||||
return event.message;
|
||||
}
|
||||
if (!snapshot) {
|
||||
throw new AnthropicError(`Unexpected event order, got ${event.type} before "message_start"`);
|
||||
}
|
||||
switch (event.type) {
|
||||
case 'message_stop':
|
||||
return snapshot;
|
||||
case 'message_delta':
|
||||
snapshot.stop_reason = event.delta.stop_reason;
|
||||
snapshot.stop_sequence = event.delta.stop_sequence;
|
||||
snapshot.usage.output_tokens = event.usage.output_tokens;
|
||||
return snapshot;
|
||||
case 'content_block_start':
|
||||
snapshot.content.push(event.content_block);
|
||||
return snapshot;
|
||||
case 'content_block_delta': {
|
||||
const snapshotContent = snapshot.content.at(event.index);
|
||||
switch (event.delta.type) {
|
||||
case 'text_delta': {
|
||||
if (snapshotContent?.type === 'text') {
|
||||
snapshotContent.text += event.delta.text;
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'citations_delta': {
|
||||
if (snapshotContent?.type === 'text') {
|
||||
snapshotContent.citations ?? (snapshotContent.citations = []);
|
||||
snapshotContent.citations.push(event.delta.citation);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'input_json_delta': {
|
||||
if (snapshotContent?.type === 'tool_use') {
|
||||
// we need to keep track of the raw JSON string as well so that we can
|
||||
// re-parse it for each delta, for now we just store it as an untyped
|
||||
// non-enumerable property on the snapshot
|
||||
let jsonBuf = snapshotContent[JSON_BUF_PROPERTY] || '';
|
||||
jsonBuf += event.delta.partial_json;
|
||||
Object.defineProperty(snapshotContent, JSON_BUF_PROPERTY, {
|
||||
value: jsonBuf,
|
||||
enumerable: false,
|
||||
writable: true,
|
||||
});
|
||||
if (jsonBuf) {
|
||||
snapshotContent.input = partialParse(jsonBuf);
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'thinking_delta': {
|
||||
if (snapshotContent?.type === 'thinking') {
|
||||
snapshotContent.thinking += event.delta.thinking;
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'signature_delta': {
|
||||
if (snapshotContent?.type === 'thinking') {
|
||||
snapshotContent.signature += event.delta.signature;
|
||||
}
|
||||
break;
|
||||
}
|
||||
default:
|
||||
checkNever(event.delta);
|
||||
}
|
||||
return snapshot;
|
||||
}
|
||||
case 'content_block_stop':
|
||||
return snapshot;
|
||||
}
|
||||
}, Symbol.asyncIterator)]() {
|
||||
const pushQueue = [];
|
||||
const readQueue = [];
|
||||
let done = false;
|
||||
this.on('streamEvent', (event) => {
|
||||
const reader = readQueue.shift();
|
||||
if (reader) {
|
||||
reader.resolve(event);
|
||||
}
|
||||
else {
|
||||
pushQueue.push(event);
|
||||
}
|
||||
});
|
||||
this.on('end', () => {
|
||||
done = true;
|
||||
for (const reader of readQueue) {
|
||||
reader.resolve(undefined);
|
||||
}
|
||||
readQueue.length = 0;
|
||||
});
|
||||
this.on('abort', (err) => {
|
||||
done = true;
|
||||
for (const reader of readQueue) {
|
||||
reader.reject(err);
|
||||
}
|
||||
readQueue.length = 0;
|
||||
});
|
||||
this.on('error', (err) => {
|
||||
done = true;
|
||||
for (const reader of readQueue) {
|
||||
reader.reject(err);
|
||||
}
|
||||
readQueue.length = 0;
|
||||
});
|
||||
return {
|
||||
next: async () => {
|
||||
if (!pushQueue.length) {
|
||||
if (done) {
|
||||
return { value: undefined, done: true };
|
||||
}
|
||||
return new Promise((resolve, reject) => readQueue.push({ resolve, reject })).then((chunk) => (chunk ? { value: chunk, done: false } : { value: undefined, done: true }));
|
||||
}
|
||||
const chunk = pushQueue.shift();
|
||||
return { value: chunk, done: false };
|
||||
},
|
||||
return: async () => {
|
||||
this.abort();
|
||||
return { value: undefined, done: true };
|
||||
},
|
||||
};
|
||||
}
|
||||
toReadableStream() {
|
||||
const stream = new Stream(this[Symbol.asyncIterator].bind(this), this.controller);
|
||||
return stream.toReadableStream();
|
||||
}
|
||||
}
|
||||
// used to ensure exhaustive case matching without throwing a runtime error
|
||||
function checkNever(x) { }
|
||||
//# sourceMappingURL=BetaMessageStream.mjs.map
|
||||
1
vendor/sdk/lib/BetaMessageStream.mjs.map
vendored
Normal file
1
vendor/sdk/lib/BetaMessageStream.mjs.map
vendored
Normal file
File diff suppressed because one or more lines are too long
114
vendor/sdk/lib/MessageStream.d.ts
vendored
Normal file
114
vendor/sdk/lib/MessageStream.d.ts
vendored
Normal file
@@ -0,0 +1,114 @@
|
||||
|
||||
import * as Core from '@anthropic-ai/sdk/core';
|
||||
import { AnthropicError, APIUserAbortError } from '@anthropic-ai/sdk/error';
|
||||
import { type ContentBlock, Messages, type Message, type MessageStreamEvent, type MessageParam, type MessageCreateParams, type MessageCreateParamsBase, type TextCitation } from '@anthropic-ai/sdk/resources/messages';
|
||||
import { type ReadableStream, type Response } from '@anthropic-ai/sdk/_shims/index';
|
||||
export interface MessageStreamEvents {
|
||||
connect: () => void;
|
||||
streamEvent: (event: MessageStreamEvent, snapshot: Message) => void;
|
||||
text: (textDelta: string, textSnapshot: string) => void;
|
||||
citation: (citation: TextCitation, citationsSnapshot: TextCitation[]) => void;
|
||||
inputJson: (partialJson: string, jsonSnapshot: unknown) => void;
|
||||
thinking: (thinkingDelta: string, thinkingSnapshot: string) => void;
|
||||
message: (message: Message) => void;
|
||||
contentBlock: (content: ContentBlock) => void;
|
||||
finalMessage: (message: Message) => void;
|
||||
error: (error: AnthropicError) => void;
|
||||
abort: (error: APIUserAbortError) => void;
|
||||
end: () => void;
|
||||
}
|
||||
export declare class MessageStream implements AsyncIterable<MessageStreamEvent> {
|
||||
#private;
|
||||
messages: MessageParam[];
|
||||
receivedMessages: Message[];
|
||||
controller: AbortController;
|
||||
constructor();
|
||||
get response(): Response | null | undefined;
|
||||
get request_id(): string | null | undefined;
|
||||
/**
|
||||
* Returns the `MessageStream` data, the raw `Response` instance and the ID of the request,
|
||||
* returned vie the `request-id` header which is useful for debugging requests and resporting
|
||||
* issues to Anthropic.
|
||||
*
|
||||
* This is the same as the `APIPromise.withResponse()` method.
|
||||
*
|
||||
* This method will raise an error if you created the stream using `MessageStream.fromReadableStream`
|
||||
* as no `Response` is available.
|
||||
*/
|
||||
withResponse(): Promise<{
|
||||
data: MessageStream;
|
||||
response: Response;
|
||||
request_id: string | null | undefined;
|
||||
}>;
|
||||
/**
|
||||
* Intended for use on the frontend, consuming a stream produced with
|
||||
* `.toReadableStream()` on the backend.
|
||||
*
|
||||
* Note that messages sent to the model do not appear in `.on('message')`
|
||||
* in this context.
|
||||
*/
|
||||
static fromReadableStream(stream: ReadableStream): MessageStream;
|
||||
static createMessage(messages: Messages, params: MessageCreateParamsBase, options?: Core.RequestOptions): MessageStream;
|
||||
protected _run(executor: () => Promise<any>): void;
|
||||
protected _addMessageParam(message: MessageParam): void;
|
||||
protected _addMessage(message: Message, emit?: boolean): void;
|
||||
protected _createMessage(messages: Messages, params: MessageCreateParams, options?: Core.RequestOptions): Promise<void>;
|
||||
protected _connected(response: Response | null): void;
|
||||
get ended(): boolean;
|
||||
get errored(): boolean;
|
||||
get aborted(): boolean;
|
||||
abort(): void;
|
||||
/**
|
||||
* Adds the listener function to the end of the listeners array for the event.
|
||||
* No checks are made to see if the listener has already been added. Multiple calls passing
|
||||
* the same combination of event and listener will result in the listener being added, and
|
||||
* called, multiple times.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
on<Event extends keyof MessageStreamEvents>(event: Event, listener: MessageStreamEvents[Event]): this;
|
||||
/**
|
||||
* Removes the specified listener from the listener array for the event.
|
||||
* off() will remove, at most, one instance of a listener from the listener array. If any single
|
||||
* listener has been added multiple times to the listener array for the specified event, then
|
||||
* off() must be called multiple times to remove each instance.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
off<Event extends keyof MessageStreamEvents>(event: Event, listener: MessageStreamEvents[Event]): this;
|
||||
/**
|
||||
* Adds a one-time listener function for the event. The next time the event is triggered,
|
||||
* this listener is removed and then invoked.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
once<Event extends keyof MessageStreamEvents>(event: Event, listener: MessageStreamEvents[Event]): this;
|
||||
/**
|
||||
* This is similar to `.once()`, but returns a Promise that resolves the next time
|
||||
* the event is triggered, instead of calling a listener callback.
|
||||
* @returns a Promise that resolves the next time given event is triggered,
|
||||
* or rejects if an error is emitted. (If you request the 'error' event,
|
||||
* returns a promise that resolves with the error).
|
||||
*
|
||||
* Example:
|
||||
*
|
||||
* const message = await stream.emitted('message') // rejects if the stream errors
|
||||
*/
|
||||
emitted<Event extends keyof MessageStreamEvents>(event: Event): Promise<Parameters<MessageStreamEvents[Event]> extends [infer Param] ? Param : Parameters<MessageStreamEvents[Event]> extends [] ? void : Parameters<MessageStreamEvents[Event]>>;
|
||||
done(): Promise<void>;
|
||||
get currentMessage(): Message | undefined;
|
||||
/**
|
||||
* @returns a promise that resolves with the the final assistant Message response,
|
||||
* or rejects if an error occurred or the stream ended prematurely without producing a Message.
|
||||
*/
|
||||
finalMessage(): Promise<Message>;
|
||||
/**
|
||||
* @returns a promise that resolves with the the final assistant Message's text response, concatenated
|
||||
* together if there are more than one text blocks.
|
||||
* Rejects if an error occurred or the stream ended prematurely without producing a Message.
|
||||
*/
|
||||
finalText(): Promise<string>;
|
||||
protected _emit<Event extends keyof MessageStreamEvents>(event: Event, ...args: Parameters<MessageStreamEvents[Event]>): void;
|
||||
protected _emitFinal(): void;
|
||||
protected _fromReadableStream(readableStream: ReadableStream, options?: Core.RequestOptions): Promise<void>;
|
||||
[Symbol.asyncIterator](): AsyncIterator<MessageStreamEvent>;
|
||||
toReadableStream(): ReadableStream;
|
||||
}
|
||||
//# sourceMappingURL=MessageStream.d.ts.map
|
||||
1
vendor/sdk/lib/MessageStream.d.ts.map
vendored
Normal file
1
vendor/sdk/lib/MessageStream.d.ts.map
vendored
Normal file
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"MessageStream.d.ts","sourceRoot":"","sources":["../src/lib/MessageStream.ts"],"names":[],"mappings":";AAAA,OAAO,KAAK,IAAI,MAAM,wBAAwB,CAAC;AAC/C,OAAO,EAAE,cAAc,EAAE,iBAAiB,EAAE,MAAM,yBAAyB,CAAC;AAC5E,OAAO,EACL,KAAK,YAAY,EACjB,QAAQ,EACR,KAAK,OAAO,EACZ,KAAK,kBAAkB,EACvB,KAAK,YAAY,EACjB,KAAK,mBAAmB,EACxB,KAAK,uBAAuB,EAE5B,KAAK,YAAY,EAClB,MAAM,sCAAsC,CAAC;AAC9C,OAAO,EAAE,KAAK,cAAc,EAAE,KAAK,QAAQ,EAAE,MAAM,gCAAgC,CAAC;AAIpF,MAAM,WAAW,mBAAmB;IAClC,OAAO,EAAE,MAAM,IAAI,CAAC;IACpB,WAAW,EAAE,CAAC,KAAK,EAAE,kBAAkB,EAAE,QAAQ,EAAE,OAAO,KAAK,IAAI,CAAC;IACpE,IAAI,EAAE,CAAC,SAAS,EAAE,MAAM,EAAE,YAAY,EAAE,MAAM,KAAK,IAAI,CAAC;IACxD,QAAQ,EAAE,CAAC,QAAQ,EAAE,YAAY,EAAE,iBAAiB,EAAE,YAAY,EAAE,KAAK,IAAI,CAAC;IAC9E,SAAS,EAAE,CAAC,WAAW,EAAE,MAAM,EAAE,YAAY,EAAE,OAAO,KAAK,IAAI,CAAC;IAChE,QAAQ,EAAE,CAAC,aAAa,EAAE,MAAM,EAAE,gBAAgB,EAAE,MAAM,KAAK,IAAI,CAAC;IACpE,OAAO,EAAE,CAAC,OAAO,EAAE,OAAO,KAAK,IAAI,CAAC;IACpC,YAAY,EAAE,CAAC,OAAO,EAAE,YAAY,KAAK,IAAI,CAAC;IAC9C,YAAY,EAAE,CAAC,OAAO,EAAE,OAAO,KAAK,IAAI,CAAC;IACzC,KAAK,EAAE,CAAC,KAAK,EAAE,cAAc,KAAK,IAAI,CAAC;IACvC,KAAK,EAAE,CAAC,KAAK,EAAE,iBAAiB,KAAK,IAAI,CAAC;IAC1C,GAAG,EAAE,MAAM,IAAI,CAAC;CACjB;AASD,qBAAa,aAAc,YAAW,aAAa,CAAC,kBAAkB,CAAC;;IACrE,QAAQ,EAAE,YAAY,EAAE,CAAM;IAC9B,gBAAgB,EAAE,OAAO,EAAE,CAAM;IAGjC,UAAU,EAAE,eAAe,CAAyB;;IAsCpD,IAAI,QAAQ,IAAI,QAAQ,GAAG,IAAI,GAAG,SAAS,CAE1C;IAED,IAAI,UAAU,IAAI,MAAM,GAAG,IAAI,GAAG,SAAS,CAE1C;IAED;;;;;;;;;OASG;IACG,YAAY,IAAI,OAAO,CAAC;QAC5B,IAAI,EAAE,aAAa,CAAC;QACpB,QAAQ,EAAE,QAAQ,CAAC;QACnB,UAAU,EAAE,MAAM,GAAG,IAAI,GAAG,SAAS,CAAC;KACvC,CAAC;IAaF;;;;;;OAMG;IACH,MAAM,CAAC,kBAAkB,CAAC,MAAM,EAAE,cAAc,GAAG,aAAa;IAMhE,MAAM,CAAC,aAAa,CAClB,QAAQ,EAAE,QAAQ,EAClB,MAAM,EAAE,uBAAuB,EAC/B,OAAO,CAAC,EAAE,IAAI,CAAC,cAAc,GAC5B,aAAa;IAehB,SAAS,CAAC,IAAI,CAAC,QAAQ,EAAE,MAAM,OAAO,CAAC,GAAG,CAAC;IAO3C,SAAS,CAAC,gBAAgB,CAAC,OAAO,EAAE,YAAY;IAIhD,SAAS,CAAC,WAAW,CAAC,OAAO,EAAE,OAAO,EAAE,IAAI,UAAO;cAOnC,cAAc,CAC5B,QAAQ,EAAE,QAAQ,EAClB,MAAM,EAAE,mBAAmB,EAC3B,OAAO,CAAC,EAAE,IAAI,CAAC,cAAc,GAC5B,OAAO,CAAC,IAAI,CAAC;IAoBhB,SAAS,CAAC,UAAU,CAAC,QAAQ,EAAE,QAAQ,GAAG,IAAI;IAQ9C,IAAI,KAAK,IAAI,OAAO,CAEnB;IAED,IAAI,OAAO,IAAI,OAAO,CAErB;IAED,IAAI,OAAO,IAAI,OAAO,CAErB;IAED,KAAK;IAIL;;;;;;OAMG;IACH,EAAE,CAAC,KAAK,SAAS,MAAM,mBAAmB,EAAE,KAAK,EAAE,KAAK,EAAE,QAAQ,EAAE,mBAAmB,CAAC,KAAK,CAAC,GAAG,IAAI;IAOrG;;;;;;OAMG;IACH,GAAG,CAAC,KAAK,SAAS,MAAM,mBAAmB,EAAE,KAAK,EAAE,KAAK,EAAE,QAAQ,EAAE,mBAAmB,CAAC,KAAK,CAAC,GAAG,IAAI;IAQtG;;;;OAIG;IACH,IAAI,CAAC,KAAK,SAAS,MAAM,mBAAmB,EAAE,KAAK,EAAE,KAAK,EAAE,QAAQ,EAAE,mBAAmB,CAAC,KAAK,CAAC,GAAG,IAAI;IAOvG;;;;;;;;;;OAUG;IACH,OAAO,CAAC,KAAK,SAAS,MAAM,mBAAmB,EAC7C,KAAK,EAAE,KAAK,GACX,OAAO,CACR,UAAU,CAAC,mBAAmB,CAAC,KAAK,CAAC,CAAC,SAAS,CAAC,MAAM,KAAK,CAAC,GAAG,KAAK,GAClE,UAAU,CAAC,mBAAmB,CAAC,KAAK,CAAC,CAAC,SAAS,EAAE,GAAG,IAAI,GACxD,UAAU,CAAC,mBAAmB,CAAC,KAAK,CAAC,CAAC,CACzC;IAQK,IAAI,IAAI,OAAO,CAAC,IAAI,CAAC;IAK3B,IAAI,cAAc,IAAI,OAAO,GAAG,SAAS,CAExC;IASD;;;OAGG;IACG,YAAY,IAAI,OAAO,CAAC,OAAO,CAAC;IAmBtC;;;;OAIG;IACG,SAAS,IAAI,OAAO,CAAC,MAAM,CAAC;IA0BlC,SAAS,CAAC,KAAK,CAAC,KAAK,SAAS,MAAM,mBAAmB,EACrD,KAAK,EAAE,KAAK,EACZ,GAAG,IAAI,EAAE,UAAU,CAAC,mBAAmB,CAAC,KAAK,CAAC,CAAC;IA8CjD,SAAS,CAAC,UAAU;cAmFJ,mBAAmB,CACjC,cAAc,EAAE,cAAc,EAC9B,OAAO,CAAC,EAAE,IAAI,CAAC,cAAc,GAC5B,OAAO,CAAC,IAAI,CAAC;IA4GhB,CAAC,MAAM,CAAC,aAAa,CAAC,IAAI,aAAa,CAAC,kBAAkB,CAAC;IA6D3D,gBAAgB,IAAI,cAAc;CAInC"}
|
||||
547
vendor/sdk/lib/MessageStream.js
vendored
Normal file
547
vendor/sdk/lib/MessageStream.js
vendored
Normal file
@@ -0,0 +1,547 @@
|
||||
"use strict";
|
||||
var __classPrivateFieldSet = (this && this.__classPrivateFieldSet) || function (receiver, state, value, kind, f) {
|
||||
if (kind === "m") throw new TypeError("Private method is not writable");
|
||||
if (kind === "a" && !f) throw new TypeError("Private accessor was defined without a setter");
|
||||
if (typeof state === "function" ? receiver !== state || !f : !state.has(receiver)) throw new TypeError("Cannot write private member to an object whose class did not declare it");
|
||||
return (kind === "a" ? f.call(receiver, value) : f ? f.value = value : state.set(receiver, value)), value;
|
||||
};
|
||||
var __classPrivateFieldGet = (this && this.__classPrivateFieldGet) || function (receiver, state, kind, f) {
|
||||
if (kind === "a" && !f) throw new TypeError("Private accessor was defined without a getter");
|
||||
if (typeof state === "function" ? receiver !== state || !f : !state.has(receiver)) throw new TypeError("Cannot read private member from an object whose class did not declare it");
|
||||
return kind === "m" ? f : kind === "a" ? f.call(receiver) : f ? f.value : state.get(receiver);
|
||||
};
|
||||
var _MessageStream_instances, _MessageStream_currentMessageSnapshot, _MessageStream_connectedPromise, _MessageStream_resolveConnectedPromise, _MessageStream_rejectConnectedPromise, _MessageStream_endPromise, _MessageStream_resolveEndPromise, _MessageStream_rejectEndPromise, _MessageStream_listeners, _MessageStream_ended, _MessageStream_errored, _MessageStream_aborted, _MessageStream_catchingPromiseCreated, _MessageStream_response, _MessageStream_request_id, _MessageStream_getFinalMessage, _MessageStream_getFinalText, _MessageStream_handleError, _MessageStream_beginRequest, _MessageStream_addStreamEvent, _MessageStream_endRequest, _MessageStream_accumulateMessage;
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
exports.MessageStream = void 0;
|
||||
const error_1 = require("@anthropic-ai/sdk/error");
|
||||
const streaming_1 = require("@anthropic-ai/sdk/streaming");
|
||||
const parser_1 = require("../_vendor/partial-json-parser/parser.js");
|
||||
const JSON_BUF_PROPERTY = '__json_buf';
|
||||
class MessageStream {
|
||||
constructor() {
|
||||
_MessageStream_instances.add(this);
|
||||
this.messages = [];
|
||||
this.receivedMessages = [];
|
||||
_MessageStream_currentMessageSnapshot.set(this, void 0);
|
||||
this.controller = new AbortController();
|
||||
_MessageStream_connectedPromise.set(this, void 0);
|
||||
_MessageStream_resolveConnectedPromise.set(this, () => { });
|
||||
_MessageStream_rejectConnectedPromise.set(this, () => { });
|
||||
_MessageStream_endPromise.set(this, void 0);
|
||||
_MessageStream_resolveEndPromise.set(this, () => { });
|
||||
_MessageStream_rejectEndPromise.set(this, () => { });
|
||||
_MessageStream_listeners.set(this, {});
|
||||
_MessageStream_ended.set(this, false);
|
||||
_MessageStream_errored.set(this, false);
|
||||
_MessageStream_aborted.set(this, false);
|
||||
_MessageStream_catchingPromiseCreated.set(this, false);
|
||||
_MessageStream_response.set(this, void 0);
|
||||
_MessageStream_request_id.set(this, void 0);
|
||||
_MessageStream_handleError.set(this, (error) => {
|
||||
__classPrivateFieldSet(this, _MessageStream_errored, true, "f");
|
||||
if (error instanceof Error && error.name === 'AbortError') {
|
||||
error = new error_1.APIUserAbortError();
|
||||
}
|
||||
if (error instanceof error_1.APIUserAbortError) {
|
||||
__classPrivateFieldSet(this, _MessageStream_aborted, true, "f");
|
||||
return this._emit('abort', error);
|
||||
}
|
||||
if (error instanceof error_1.AnthropicError) {
|
||||
return this._emit('error', error);
|
||||
}
|
||||
if (error instanceof Error) {
|
||||
const anthropicError = new error_1.AnthropicError(error.message);
|
||||
// @ts-ignore
|
||||
anthropicError.cause = error;
|
||||
return this._emit('error', anthropicError);
|
||||
}
|
||||
return this._emit('error', new error_1.AnthropicError(String(error)));
|
||||
});
|
||||
__classPrivateFieldSet(this, _MessageStream_connectedPromise, new Promise((resolve, reject) => {
|
||||
__classPrivateFieldSet(this, _MessageStream_resolveConnectedPromise, resolve, "f");
|
||||
__classPrivateFieldSet(this, _MessageStream_rejectConnectedPromise, reject, "f");
|
||||
}), "f");
|
||||
__classPrivateFieldSet(this, _MessageStream_endPromise, new Promise((resolve, reject) => {
|
||||
__classPrivateFieldSet(this, _MessageStream_resolveEndPromise, resolve, "f");
|
||||
__classPrivateFieldSet(this, _MessageStream_rejectEndPromise, reject, "f");
|
||||
}), "f");
|
||||
// Don't let these promises cause unhandled rejection errors.
|
||||
// we will manually cause an unhandled rejection error later
|
||||
// if the user hasn't registered any error listener or called
|
||||
// any promise-returning method.
|
||||
__classPrivateFieldGet(this, _MessageStream_connectedPromise, "f").catch(() => { });
|
||||
__classPrivateFieldGet(this, _MessageStream_endPromise, "f").catch(() => { });
|
||||
}
|
||||
get response() {
|
||||
return __classPrivateFieldGet(this, _MessageStream_response, "f");
|
||||
}
|
||||
get request_id() {
|
||||
return __classPrivateFieldGet(this, _MessageStream_request_id, "f");
|
||||
}
|
||||
/**
|
||||
* Returns the `MessageStream` data, the raw `Response` instance and the ID of the request,
|
||||
* returned vie the `request-id` header which is useful for debugging requests and resporting
|
||||
* issues to Anthropic.
|
||||
*
|
||||
* This is the same as the `APIPromise.withResponse()` method.
|
||||
*
|
||||
* This method will raise an error if you created the stream using `MessageStream.fromReadableStream`
|
||||
* as no `Response` is available.
|
||||
*/
|
||||
async withResponse() {
|
||||
const response = await __classPrivateFieldGet(this, _MessageStream_connectedPromise, "f");
|
||||
if (!response) {
|
||||
throw new Error('Could not resolve a `Response` object');
|
||||
}
|
||||
return {
|
||||
data: this,
|
||||
response,
|
||||
request_id: response.headers.get('request-id'),
|
||||
};
|
||||
}
|
||||
/**
|
||||
* Intended for use on the frontend, consuming a stream produced with
|
||||
* `.toReadableStream()` on the backend.
|
||||
*
|
||||
* Note that messages sent to the model do not appear in `.on('message')`
|
||||
* in this context.
|
||||
*/
|
||||
static fromReadableStream(stream) {
|
||||
const runner = new MessageStream();
|
||||
runner._run(() => runner._fromReadableStream(stream));
|
||||
return runner;
|
||||
}
|
||||
static createMessage(messages, params, options) {
|
||||
const runner = new MessageStream();
|
||||
for (const message of params.messages) {
|
||||
runner._addMessageParam(message);
|
||||
}
|
||||
runner._run(() => runner._createMessage(messages, { ...params, stream: true }, { ...options, headers: { ...options?.headers, 'X-Stainless-Helper-Method': 'stream' } }));
|
||||
return runner;
|
||||
}
|
||||
_run(executor) {
|
||||
executor().then(() => {
|
||||
this._emitFinal();
|
||||
this._emit('end');
|
||||
}, __classPrivateFieldGet(this, _MessageStream_handleError, "f"));
|
||||
}
|
||||
_addMessageParam(message) {
|
||||
this.messages.push(message);
|
||||
}
|
||||
_addMessage(message, emit = true) {
|
||||
this.receivedMessages.push(message);
|
||||
if (emit) {
|
||||
this._emit('message', message);
|
||||
}
|
||||
}
|
||||
async _createMessage(messages, params, options) {
|
||||
const signal = options?.signal;
|
||||
if (signal) {
|
||||
if (signal.aborted)
|
||||
this.controller.abort();
|
||||
signal.addEventListener('abort', () => this.controller.abort());
|
||||
}
|
||||
__classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_beginRequest).call(this);
|
||||
const { response, data: stream } = await messages
|
||||
.create({ ...params, stream: true }, { ...options, signal: this.controller.signal })
|
||||
.withResponse();
|
||||
this._connected(response);
|
||||
for await (const event of stream) {
|
||||
__classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_addStreamEvent).call(this, event);
|
||||
}
|
||||
if (stream.controller.signal?.aborted) {
|
||||
throw new error_1.APIUserAbortError();
|
||||
}
|
||||
__classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_endRequest).call(this);
|
||||
}
|
||||
_connected(response) {
|
||||
if (this.ended)
|
||||
return;
|
||||
__classPrivateFieldSet(this, _MessageStream_response, response, "f");
|
||||
__classPrivateFieldSet(this, _MessageStream_request_id, response?.headers.get('request-id'), "f");
|
||||
__classPrivateFieldGet(this, _MessageStream_resolveConnectedPromise, "f").call(this, response);
|
||||
this._emit('connect');
|
||||
}
|
||||
get ended() {
|
||||
return __classPrivateFieldGet(this, _MessageStream_ended, "f");
|
||||
}
|
||||
get errored() {
|
||||
return __classPrivateFieldGet(this, _MessageStream_errored, "f");
|
||||
}
|
||||
get aborted() {
|
||||
return __classPrivateFieldGet(this, _MessageStream_aborted, "f");
|
||||
}
|
||||
abort() {
|
||||
this.controller.abort();
|
||||
}
|
||||
/**
|
||||
* Adds the listener function to the end of the listeners array for the event.
|
||||
* No checks are made to see if the listener has already been added. Multiple calls passing
|
||||
* the same combination of event and listener will result in the listener being added, and
|
||||
* called, multiple times.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
on(event, listener) {
|
||||
const listeners = __classPrivateFieldGet(this, _MessageStream_listeners, "f")[event] || (__classPrivateFieldGet(this, _MessageStream_listeners, "f")[event] = []);
|
||||
listeners.push({ listener });
|
||||
return this;
|
||||
}
|
||||
/**
|
||||
* Removes the specified listener from the listener array for the event.
|
||||
* off() will remove, at most, one instance of a listener from the listener array. If any single
|
||||
* listener has been added multiple times to the listener array for the specified event, then
|
||||
* off() must be called multiple times to remove each instance.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
off(event, listener) {
|
||||
const listeners = __classPrivateFieldGet(this, _MessageStream_listeners, "f")[event];
|
||||
if (!listeners)
|
||||
return this;
|
||||
const index = listeners.findIndex((l) => l.listener === listener);
|
||||
if (index >= 0)
|
||||
listeners.splice(index, 1);
|
||||
return this;
|
||||
}
|
||||
/**
|
||||
* Adds a one-time listener function for the event. The next time the event is triggered,
|
||||
* this listener is removed and then invoked.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
once(event, listener) {
|
||||
const listeners = __classPrivateFieldGet(this, _MessageStream_listeners, "f")[event] || (__classPrivateFieldGet(this, _MessageStream_listeners, "f")[event] = []);
|
||||
listeners.push({ listener, once: true });
|
||||
return this;
|
||||
}
|
||||
/**
|
||||
* This is similar to `.once()`, but returns a Promise that resolves the next time
|
||||
* the event is triggered, instead of calling a listener callback.
|
||||
* @returns a Promise that resolves the next time given event is triggered,
|
||||
* or rejects if an error is emitted. (If you request the 'error' event,
|
||||
* returns a promise that resolves with the error).
|
||||
*
|
||||
* Example:
|
||||
*
|
||||
* const message = await stream.emitted('message') // rejects if the stream errors
|
||||
*/
|
||||
emitted(event) {
|
||||
return new Promise((resolve, reject) => {
|
||||
__classPrivateFieldSet(this, _MessageStream_catchingPromiseCreated, true, "f");
|
||||
if (event !== 'error')
|
||||
this.once('error', reject);
|
||||
this.once(event, resolve);
|
||||
});
|
||||
}
|
||||
async done() {
|
||||
__classPrivateFieldSet(this, _MessageStream_catchingPromiseCreated, true, "f");
|
||||
await __classPrivateFieldGet(this, _MessageStream_endPromise, "f");
|
||||
}
|
||||
get currentMessage() {
|
||||
return __classPrivateFieldGet(this, _MessageStream_currentMessageSnapshot, "f");
|
||||
}
|
||||
/**
|
||||
* @returns a promise that resolves with the the final assistant Message response,
|
||||
* or rejects if an error occurred or the stream ended prematurely without producing a Message.
|
||||
*/
|
||||
async finalMessage() {
|
||||
await this.done();
|
||||
return __classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_getFinalMessage).call(this);
|
||||
}
|
||||
/**
|
||||
* @returns a promise that resolves with the the final assistant Message's text response, concatenated
|
||||
* together if there are more than one text blocks.
|
||||
* Rejects if an error occurred or the stream ended prematurely without producing a Message.
|
||||
*/
|
||||
async finalText() {
|
||||
await this.done();
|
||||
return __classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_getFinalText).call(this);
|
||||
}
|
||||
_emit(event, ...args) {
|
||||
// make sure we don't emit any MessageStreamEvents after end
|
||||
if (__classPrivateFieldGet(this, _MessageStream_ended, "f"))
|
||||
return;
|
||||
if (event === 'end') {
|
||||
__classPrivateFieldSet(this, _MessageStream_ended, true, "f");
|
||||
__classPrivateFieldGet(this, _MessageStream_resolveEndPromise, "f").call(this);
|
||||
}
|
||||
const listeners = __classPrivateFieldGet(this, _MessageStream_listeners, "f")[event];
|
||||
if (listeners) {
|
||||
__classPrivateFieldGet(this, _MessageStream_listeners, "f")[event] = listeners.filter((l) => !l.once);
|
||||
listeners.forEach(({ listener }) => listener(...args));
|
||||
}
|
||||
if (event === 'abort') {
|
||||
const error = args[0];
|
||||
if (!__classPrivateFieldGet(this, _MessageStream_catchingPromiseCreated, "f") && !listeners?.length) {
|
||||
Promise.reject(error);
|
||||
}
|
||||
__classPrivateFieldGet(this, _MessageStream_rejectConnectedPromise, "f").call(this, error);
|
||||
__classPrivateFieldGet(this, _MessageStream_rejectEndPromise, "f").call(this, error);
|
||||
this._emit('end');
|
||||
return;
|
||||
}
|
||||
if (event === 'error') {
|
||||
// NOTE: _emit('error', error) should only be called from #handleError().
|
||||
const error = args[0];
|
||||
if (!__classPrivateFieldGet(this, _MessageStream_catchingPromiseCreated, "f") && !listeners?.length) {
|
||||
// Trigger an unhandled rejection if the user hasn't registered any error handlers.
|
||||
// If you are seeing stack traces here, make sure to handle errors via either:
|
||||
// - runner.on('error', () => ...)
|
||||
// - await runner.done()
|
||||
// - await runner.final...()
|
||||
// - etc.
|
||||
Promise.reject(error);
|
||||
}
|
||||
__classPrivateFieldGet(this, _MessageStream_rejectConnectedPromise, "f").call(this, error);
|
||||
__classPrivateFieldGet(this, _MessageStream_rejectEndPromise, "f").call(this, error);
|
||||
this._emit('end');
|
||||
}
|
||||
}
|
||||
_emitFinal() {
|
||||
const finalMessage = this.receivedMessages.at(-1);
|
||||
if (finalMessage) {
|
||||
this._emit('finalMessage', __classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_getFinalMessage).call(this));
|
||||
}
|
||||
}
|
||||
async _fromReadableStream(readableStream, options) {
|
||||
const signal = options?.signal;
|
||||
if (signal) {
|
||||
if (signal.aborted)
|
||||
this.controller.abort();
|
||||
signal.addEventListener('abort', () => this.controller.abort());
|
||||
}
|
||||
__classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_beginRequest).call(this);
|
||||
this._connected(null);
|
||||
const stream = streaming_1.Stream.fromReadableStream(readableStream, this.controller);
|
||||
for await (const event of stream) {
|
||||
__classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_addStreamEvent).call(this, event);
|
||||
}
|
||||
if (stream.controller.signal?.aborted) {
|
||||
throw new error_1.APIUserAbortError();
|
||||
}
|
||||
__classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_endRequest).call(this);
|
||||
}
|
||||
[(_MessageStream_currentMessageSnapshot = new WeakMap(), _MessageStream_connectedPromise = new WeakMap(), _MessageStream_resolveConnectedPromise = new WeakMap(), _MessageStream_rejectConnectedPromise = new WeakMap(), _MessageStream_endPromise = new WeakMap(), _MessageStream_resolveEndPromise = new WeakMap(), _MessageStream_rejectEndPromise = new WeakMap(), _MessageStream_listeners = new WeakMap(), _MessageStream_ended = new WeakMap(), _MessageStream_errored = new WeakMap(), _MessageStream_aborted = new WeakMap(), _MessageStream_catchingPromiseCreated = new WeakMap(), _MessageStream_response = new WeakMap(), _MessageStream_request_id = new WeakMap(), _MessageStream_handleError = new WeakMap(), _MessageStream_instances = new WeakSet(), _MessageStream_getFinalMessage = function _MessageStream_getFinalMessage() {
|
||||
if (this.receivedMessages.length === 0) {
|
||||
throw new error_1.AnthropicError('stream ended without producing a Message with role=assistant');
|
||||
}
|
||||
return this.receivedMessages.at(-1);
|
||||
}, _MessageStream_getFinalText = function _MessageStream_getFinalText() {
|
||||
if (this.receivedMessages.length === 0) {
|
||||
throw new error_1.AnthropicError('stream ended without producing a Message with role=assistant');
|
||||
}
|
||||
const textBlocks = this.receivedMessages
|
||||
.at(-1)
|
||||
.content.filter((block) => block.type === 'text')
|
||||
.map((block) => block.text);
|
||||
if (textBlocks.length === 0) {
|
||||
throw new error_1.AnthropicError('stream ended without producing a content block with type=text');
|
||||
}
|
||||
return textBlocks.join(' ');
|
||||
}, _MessageStream_beginRequest = function _MessageStream_beginRequest() {
|
||||
if (this.ended)
|
||||
return;
|
||||
__classPrivateFieldSet(this, _MessageStream_currentMessageSnapshot, undefined, "f");
|
||||
}, _MessageStream_addStreamEvent = function _MessageStream_addStreamEvent(event) {
|
||||
if (this.ended)
|
||||
return;
|
||||
const messageSnapshot = __classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_accumulateMessage).call(this, event);
|
||||
this._emit('streamEvent', event, messageSnapshot);
|
||||
switch (event.type) {
|
||||
case 'content_block_delta': {
|
||||
const content = messageSnapshot.content.at(-1);
|
||||
switch (event.delta.type) {
|
||||
case 'text_delta': {
|
||||
if (content.type === 'text') {
|
||||
this._emit('text', event.delta.text, content.text || '');
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'citations_delta': {
|
||||
if (content.type === 'text') {
|
||||
this._emit('citation', event.delta.citation, content.citations ?? []);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'input_json_delta': {
|
||||
if (content.type === 'tool_use' && content.input) {
|
||||
this._emit('inputJson', event.delta.partial_json, content.input);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'thinking_delta': {
|
||||
if (content.type === 'thinking') {
|
||||
this._emit('thinking', event.delta.thinking, content.thinking);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'signature_delta': {
|
||||
// we don't emit anything special in this case.
|
||||
break;
|
||||
}
|
||||
default:
|
||||
checkNever(event.delta);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'message_stop': {
|
||||
this._addMessageParam(messageSnapshot);
|
||||
this._addMessage(messageSnapshot, true);
|
||||
break;
|
||||
}
|
||||
case 'content_block_stop': {
|
||||
this._emit('contentBlock', messageSnapshot.content.at(-1));
|
||||
break;
|
||||
}
|
||||
case 'message_start': {
|
||||
__classPrivateFieldSet(this, _MessageStream_currentMessageSnapshot, messageSnapshot, "f");
|
||||
break;
|
||||
}
|
||||
case 'content_block_start':
|
||||
case 'message_delta':
|
||||
break;
|
||||
}
|
||||
}, _MessageStream_endRequest = function _MessageStream_endRequest() {
|
||||
if (this.ended) {
|
||||
throw new error_1.AnthropicError(`stream has ended, this shouldn't happen`);
|
||||
}
|
||||
const snapshot = __classPrivateFieldGet(this, _MessageStream_currentMessageSnapshot, "f");
|
||||
if (!snapshot) {
|
||||
throw new error_1.AnthropicError(`request ended without sending any chunks`);
|
||||
}
|
||||
__classPrivateFieldSet(this, _MessageStream_currentMessageSnapshot, undefined, "f");
|
||||
return snapshot;
|
||||
}, _MessageStream_accumulateMessage = function _MessageStream_accumulateMessage(event) {
|
||||
let snapshot = __classPrivateFieldGet(this, _MessageStream_currentMessageSnapshot, "f");
|
||||
if (event.type === 'message_start') {
|
||||
if (snapshot) {
|
||||
throw new error_1.AnthropicError(`Unexpected event order, got ${event.type} before receiving "message_stop"`);
|
||||
}
|
||||
return event.message;
|
||||
}
|
||||
if (!snapshot) {
|
||||
throw new error_1.AnthropicError(`Unexpected event order, got ${event.type} before "message_start"`);
|
||||
}
|
||||
switch (event.type) {
|
||||
case 'message_stop':
|
||||
return snapshot;
|
||||
case 'message_delta':
|
||||
snapshot.stop_reason = event.delta.stop_reason;
|
||||
snapshot.stop_sequence = event.delta.stop_sequence;
|
||||
snapshot.usage.output_tokens = event.usage.output_tokens;
|
||||
return snapshot;
|
||||
case 'content_block_start':
|
||||
snapshot.content.push(event.content_block);
|
||||
return snapshot;
|
||||
case 'content_block_delta': {
|
||||
const snapshotContent = snapshot.content.at(event.index);
|
||||
switch (event.delta.type) {
|
||||
case 'text_delta': {
|
||||
if (snapshotContent?.type === 'text') {
|
||||
snapshotContent.text += event.delta.text;
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'citations_delta': {
|
||||
if (snapshotContent?.type === 'text') {
|
||||
snapshotContent.citations ?? (snapshotContent.citations = []);
|
||||
snapshotContent.citations.push(event.delta.citation);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'input_json_delta': {
|
||||
if (snapshotContent?.type === 'tool_use') {
|
||||
// we need to keep track of the raw JSON string as well so that we can
|
||||
// re-parse it for each delta, for now we just store it as an untyped
|
||||
// non-enumerable property on the snapshot
|
||||
let jsonBuf = snapshotContent[JSON_BUF_PROPERTY] || '';
|
||||
jsonBuf += event.delta.partial_json;
|
||||
Object.defineProperty(snapshotContent, JSON_BUF_PROPERTY, {
|
||||
value: jsonBuf,
|
||||
enumerable: false,
|
||||
writable: true,
|
||||
});
|
||||
if (jsonBuf) {
|
||||
snapshotContent.input = (0, parser_1.partialParse)(jsonBuf);
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'thinking_delta': {
|
||||
if (snapshotContent?.type === 'thinking') {
|
||||
snapshotContent.thinking += event.delta.thinking;
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'signature_delta': {
|
||||
if (snapshotContent?.type === 'thinking') {
|
||||
snapshotContent.signature += event.delta.signature;
|
||||
}
|
||||
break;
|
||||
}
|
||||
default:
|
||||
checkNever(event.delta);
|
||||
}
|
||||
return snapshot;
|
||||
}
|
||||
case 'content_block_stop':
|
||||
return snapshot;
|
||||
}
|
||||
}, Symbol.asyncIterator)]() {
|
||||
const pushQueue = [];
|
||||
const readQueue = [];
|
||||
let done = false;
|
||||
this.on('streamEvent', (event) => {
|
||||
const reader = readQueue.shift();
|
||||
if (reader) {
|
||||
reader.resolve(event);
|
||||
}
|
||||
else {
|
||||
pushQueue.push(event);
|
||||
}
|
||||
});
|
||||
this.on('end', () => {
|
||||
done = true;
|
||||
for (const reader of readQueue) {
|
||||
reader.resolve(undefined);
|
||||
}
|
||||
readQueue.length = 0;
|
||||
});
|
||||
this.on('abort', (err) => {
|
||||
done = true;
|
||||
for (const reader of readQueue) {
|
||||
reader.reject(err);
|
||||
}
|
||||
readQueue.length = 0;
|
||||
});
|
||||
this.on('error', (err) => {
|
||||
done = true;
|
||||
for (const reader of readQueue) {
|
||||
reader.reject(err);
|
||||
}
|
||||
readQueue.length = 0;
|
||||
});
|
||||
return {
|
||||
next: async () => {
|
||||
if (!pushQueue.length) {
|
||||
if (done) {
|
||||
return { value: undefined, done: true };
|
||||
}
|
||||
return new Promise((resolve, reject) => readQueue.push({ resolve, reject })).then((chunk) => (chunk ? { value: chunk, done: false } : { value: undefined, done: true }));
|
||||
}
|
||||
const chunk = pushQueue.shift();
|
||||
return { value: chunk, done: false };
|
||||
},
|
||||
return: async () => {
|
||||
this.abort();
|
||||
return { value: undefined, done: true };
|
||||
},
|
||||
};
|
||||
}
|
||||
toReadableStream() {
|
||||
const stream = new streaming_1.Stream(this[Symbol.asyncIterator].bind(this), this.controller);
|
||||
return stream.toReadableStream();
|
||||
}
|
||||
}
|
||||
exports.MessageStream = MessageStream;
|
||||
// used to ensure exhaustive case matching without throwing a runtime error
|
||||
function checkNever(x) { }
|
||||
//# sourceMappingURL=MessageStream.js.map
|
||||
1
vendor/sdk/lib/MessageStream.js.map
vendored
Normal file
1
vendor/sdk/lib/MessageStream.js.map
vendored
Normal file
File diff suppressed because one or more lines are too long
543
vendor/sdk/lib/MessageStream.mjs
vendored
Normal file
543
vendor/sdk/lib/MessageStream.mjs
vendored
Normal file
@@ -0,0 +1,543 @@
|
||||
var __classPrivateFieldSet = (this && this.__classPrivateFieldSet) || function (receiver, state, value, kind, f) {
|
||||
if (kind === "m") throw new TypeError("Private method is not writable");
|
||||
if (kind === "a" && !f) throw new TypeError("Private accessor was defined without a setter");
|
||||
if (typeof state === "function" ? receiver !== state || !f : !state.has(receiver)) throw new TypeError("Cannot write private member to an object whose class did not declare it");
|
||||
return (kind === "a" ? f.call(receiver, value) : f ? f.value = value : state.set(receiver, value)), value;
|
||||
};
|
||||
var __classPrivateFieldGet = (this && this.__classPrivateFieldGet) || function (receiver, state, kind, f) {
|
||||
if (kind === "a" && !f) throw new TypeError("Private accessor was defined without a getter");
|
||||
if (typeof state === "function" ? receiver !== state || !f : !state.has(receiver)) throw new TypeError("Cannot read private member from an object whose class did not declare it");
|
||||
return kind === "m" ? f : kind === "a" ? f.call(receiver) : f ? f.value : state.get(receiver);
|
||||
};
|
||||
var _MessageStream_instances, _MessageStream_currentMessageSnapshot, _MessageStream_connectedPromise, _MessageStream_resolveConnectedPromise, _MessageStream_rejectConnectedPromise, _MessageStream_endPromise, _MessageStream_resolveEndPromise, _MessageStream_rejectEndPromise, _MessageStream_listeners, _MessageStream_ended, _MessageStream_errored, _MessageStream_aborted, _MessageStream_catchingPromiseCreated, _MessageStream_response, _MessageStream_request_id, _MessageStream_getFinalMessage, _MessageStream_getFinalText, _MessageStream_handleError, _MessageStream_beginRequest, _MessageStream_addStreamEvent, _MessageStream_endRequest, _MessageStream_accumulateMessage;
|
||||
import { AnthropicError, APIUserAbortError } from '@anthropic-ai/sdk/error';
|
||||
import { Stream } from '@anthropic-ai/sdk/streaming';
|
||||
import { partialParse } from "../_vendor/partial-json-parser/parser.mjs";
|
||||
const JSON_BUF_PROPERTY = '__json_buf';
|
||||
export class MessageStream {
|
||||
constructor() {
|
||||
_MessageStream_instances.add(this);
|
||||
this.messages = [];
|
||||
this.receivedMessages = [];
|
||||
_MessageStream_currentMessageSnapshot.set(this, void 0);
|
||||
this.controller = new AbortController();
|
||||
_MessageStream_connectedPromise.set(this, void 0);
|
||||
_MessageStream_resolveConnectedPromise.set(this, () => { });
|
||||
_MessageStream_rejectConnectedPromise.set(this, () => { });
|
||||
_MessageStream_endPromise.set(this, void 0);
|
||||
_MessageStream_resolveEndPromise.set(this, () => { });
|
||||
_MessageStream_rejectEndPromise.set(this, () => { });
|
||||
_MessageStream_listeners.set(this, {});
|
||||
_MessageStream_ended.set(this, false);
|
||||
_MessageStream_errored.set(this, false);
|
||||
_MessageStream_aborted.set(this, false);
|
||||
_MessageStream_catchingPromiseCreated.set(this, false);
|
||||
_MessageStream_response.set(this, void 0);
|
||||
_MessageStream_request_id.set(this, void 0);
|
||||
_MessageStream_handleError.set(this, (error) => {
|
||||
__classPrivateFieldSet(this, _MessageStream_errored, true, "f");
|
||||
if (error instanceof Error && error.name === 'AbortError') {
|
||||
error = new APIUserAbortError();
|
||||
}
|
||||
if (error instanceof APIUserAbortError) {
|
||||
__classPrivateFieldSet(this, _MessageStream_aborted, true, "f");
|
||||
return this._emit('abort', error);
|
||||
}
|
||||
if (error instanceof AnthropicError) {
|
||||
return this._emit('error', error);
|
||||
}
|
||||
if (error instanceof Error) {
|
||||
const anthropicError = new AnthropicError(error.message);
|
||||
// @ts-ignore
|
||||
anthropicError.cause = error;
|
||||
return this._emit('error', anthropicError);
|
||||
}
|
||||
return this._emit('error', new AnthropicError(String(error)));
|
||||
});
|
||||
__classPrivateFieldSet(this, _MessageStream_connectedPromise, new Promise((resolve, reject) => {
|
||||
__classPrivateFieldSet(this, _MessageStream_resolveConnectedPromise, resolve, "f");
|
||||
__classPrivateFieldSet(this, _MessageStream_rejectConnectedPromise, reject, "f");
|
||||
}), "f");
|
||||
__classPrivateFieldSet(this, _MessageStream_endPromise, new Promise((resolve, reject) => {
|
||||
__classPrivateFieldSet(this, _MessageStream_resolveEndPromise, resolve, "f");
|
||||
__classPrivateFieldSet(this, _MessageStream_rejectEndPromise, reject, "f");
|
||||
}), "f");
|
||||
// Don't let these promises cause unhandled rejection errors.
|
||||
// we will manually cause an unhandled rejection error later
|
||||
// if the user hasn't registered any error listener or called
|
||||
// any promise-returning method.
|
||||
__classPrivateFieldGet(this, _MessageStream_connectedPromise, "f").catch(() => { });
|
||||
__classPrivateFieldGet(this, _MessageStream_endPromise, "f").catch(() => { });
|
||||
}
|
||||
get response() {
|
||||
return __classPrivateFieldGet(this, _MessageStream_response, "f");
|
||||
}
|
||||
get request_id() {
|
||||
return __classPrivateFieldGet(this, _MessageStream_request_id, "f");
|
||||
}
|
||||
/**
|
||||
* Returns the `MessageStream` data, the raw `Response` instance and the ID of the request,
|
||||
* returned vie the `request-id` header which is useful for debugging requests and resporting
|
||||
* issues to Anthropic.
|
||||
*
|
||||
* This is the same as the `APIPromise.withResponse()` method.
|
||||
*
|
||||
* This method will raise an error if you created the stream using `MessageStream.fromReadableStream`
|
||||
* as no `Response` is available.
|
||||
*/
|
||||
async withResponse() {
|
||||
const response = await __classPrivateFieldGet(this, _MessageStream_connectedPromise, "f");
|
||||
if (!response) {
|
||||
throw new Error('Could not resolve a `Response` object');
|
||||
}
|
||||
return {
|
||||
data: this,
|
||||
response,
|
||||
request_id: response.headers.get('request-id'),
|
||||
};
|
||||
}
|
||||
/**
|
||||
* Intended for use on the frontend, consuming a stream produced with
|
||||
* `.toReadableStream()` on the backend.
|
||||
*
|
||||
* Note that messages sent to the model do not appear in `.on('message')`
|
||||
* in this context.
|
||||
*/
|
||||
static fromReadableStream(stream) {
|
||||
const runner = new MessageStream();
|
||||
runner._run(() => runner._fromReadableStream(stream));
|
||||
return runner;
|
||||
}
|
||||
static createMessage(messages, params, options) {
|
||||
const runner = new MessageStream();
|
||||
for (const message of params.messages) {
|
||||
runner._addMessageParam(message);
|
||||
}
|
||||
runner._run(() => runner._createMessage(messages, { ...params, stream: true }, { ...options, headers: { ...options?.headers, 'X-Stainless-Helper-Method': 'stream' } }));
|
||||
return runner;
|
||||
}
|
||||
_run(executor) {
|
||||
executor().then(() => {
|
||||
this._emitFinal();
|
||||
this._emit('end');
|
||||
}, __classPrivateFieldGet(this, _MessageStream_handleError, "f"));
|
||||
}
|
||||
_addMessageParam(message) {
|
||||
this.messages.push(message);
|
||||
}
|
||||
_addMessage(message, emit = true) {
|
||||
this.receivedMessages.push(message);
|
||||
if (emit) {
|
||||
this._emit('message', message);
|
||||
}
|
||||
}
|
||||
async _createMessage(messages, params, options) {
|
||||
const signal = options?.signal;
|
||||
if (signal) {
|
||||
if (signal.aborted)
|
||||
this.controller.abort();
|
||||
signal.addEventListener('abort', () => this.controller.abort());
|
||||
}
|
||||
__classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_beginRequest).call(this);
|
||||
const { response, data: stream } = await messages
|
||||
.create({ ...params, stream: true }, { ...options, signal: this.controller.signal })
|
||||
.withResponse();
|
||||
this._connected(response);
|
||||
for await (const event of stream) {
|
||||
__classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_addStreamEvent).call(this, event);
|
||||
}
|
||||
if (stream.controller.signal?.aborted) {
|
||||
throw new APIUserAbortError();
|
||||
}
|
||||
__classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_endRequest).call(this);
|
||||
}
|
||||
_connected(response) {
|
||||
if (this.ended)
|
||||
return;
|
||||
__classPrivateFieldSet(this, _MessageStream_response, response, "f");
|
||||
__classPrivateFieldSet(this, _MessageStream_request_id, response?.headers.get('request-id'), "f");
|
||||
__classPrivateFieldGet(this, _MessageStream_resolveConnectedPromise, "f").call(this, response);
|
||||
this._emit('connect');
|
||||
}
|
||||
get ended() {
|
||||
return __classPrivateFieldGet(this, _MessageStream_ended, "f");
|
||||
}
|
||||
get errored() {
|
||||
return __classPrivateFieldGet(this, _MessageStream_errored, "f");
|
||||
}
|
||||
get aborted() {
|
||||
return __classPrivateFieldGet(this, _MessageStream_aborted, "f");
|
||||
}
|
||||
abort() {
|
||||
this.controller.abort();
|
||||
}
|
||||
/**
|
||||
* Adds the listener function to the end of the listeners array for the event.
|
||||
* No checks are made to see if the listener has already been added. Multiple calls passing
|
||||
* the same combination of event and listener will result in the listener being added, and
|
||||
* called, multiple times.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
on(event, listener) {
|
||||
const listeners = __classPrivateFieldGet(this, _MessageStream_listeners, "f")[event] || (__classPrivateFieldGet(this, _MessageStream_listeners, "f")[event] = []);
|
||||
listeners.push({ listener });
|
||||
return this;
|
||||
}
|
||||
/**
|
||||
* Removes the specified listener from the listener array for the event.
|
||||
* off() will remove, at most, one instance of a listener from the listener array. If any single
|
||||
* listener has been added multiple times to the listener array for the specified event, then
|
||||
* off() must be called multiple times to remove each instance.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
off(event, listener) {
|
||||
const listeners = __classPrivateFieldGet(this, _MessageStream_listeners, "f")[event];
|
||||
if (!listeners)
|
||||
return this;
|
||||
const index = listeners.findIndex((l) => l.listener === listener);
|
||||
if (index >= 0)
|
||||
listeners.splice(index, 1);
|
||||
return this;
|
||||
}
|
||||
/**
|
||||
* Adds a one-time listener function for the event. The next time the event is triggered,
|
||||
* this listener is removed and then invoked.
|
||||
* @returns this MessageStream, so that calls can be chained
|
||||
*/
|
||||
once(event, listener) {
|
||||
const listeners = __classPrivateFieldGet(this, _MessageStream_listeners, "f")[event] || (__classPrivateFieldGet(this, _MessageStream_listeners, "f")[event] = []);
|
||||
listeners.push({ listener, once: true });
|
||||
return this;
|
||||
}
|
||||
/**
|
||||
* This is similar to `.once()`, but returns a Promise that resolves the next time
|
||||
* the event is triggered, instead of calling a listener callback.
|
||||
* @returns a Promise that resolves the next time given event is triggered,
|
||||
* or rejects if an error is emitted. (If you request the 'error' event,
|
||||
* returns a promise that resolves with the error).
|
||||
*
|
||||
* Example:
|
||||
*
|
||||
* const message = await stream.emitted('message') // rejects if the stream errors
|
||||
*/
|
||||
emitted(event) {
|
||||
return new Promise((resolve, reject) => {
|
||||
__classPrivateFieldSet(this, _MessageStream_catchingPromiseCreated, true, "f");
|
||||
if (event !== 'error')
|
||||
this.once('error', reject);
|
||||
this.once(event, resolve);
|
||||
});
|
||||
}
|
||||
async done() {
|
||||
__classPrivateFieldSet(this, _MessageStream_catchingPromiseCreated, true, "f");
|
||||
await __classPrivateFieldGet(this, _MessageStream_endPromise, "f");
|
||||
}
|
||||
get currentMessage() {
|
||||
return __classPrivateFieldGet(this, _MessageStream_currentMessageSnapshot, "f");
|
||||
}
|
||||
/**
|
||||
* @returns a promise that resolves with the the final assistant Message response,
|
||||
* or rejects if an error occurred or the stream ended prematurely without producing a Message.
|
||||
*/
|
||||
async finalMessage() {
|
||||
await this.done();
|
||||
return __classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_getFinalMessage).call(this);
|
||||
}
|
||||
/**
|
||||
* @returns a promise that resolves with the the final assistant Message's text response, concatenated
|
||||
* together if there are more than one text blocks.
|
||||
* Rejects if an error occurred or the stream ended prematurely without producing a Message.
|
||||
*/
|
||||
async finalText() {
|
||||
await this.done();
|
||||
return __classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_getFinalText).call(this);
|
||||
}
|
||||
_emit(event, ...args) {
|
||||
// make sure we don't emit any MessageStreamEvents after end
|
||||
if (__classPrivateFieldGet(this, _MessageStream_ended, "f"))
|
||||
return;
|
||||
if (event === 'end') {
|
||||
__classPrivateFieldSet(this, _MessageStream_ended, true, "f");
|
||||
__classPrivateFieldGet(this, _MessageStream_resolveEndPromise, "f").call(this);
|
||||
}
|
||||
const listeners = __classPrivateFieldGet(this, _MessageStream_listeners, "f")[event];
|
||||
if (listeners) {
|
||||
__classPrivateFieldGet(this, _MessageStream_listeners, "f")[event] = listeners.filter((l) => !l.once);
|
||||
listeners.forEach(({ listener }) => listener(...args));
|
||||
}
|
||||
if (event === 'abort') {
|
||||
const error = args[0];
|
||||
if (!__classPrivateFieldGet(this, _MessageStream_catchingPromiseCreated, "f") && !listeners?.length) {
|
||||
Promise.reject(error);
|
||||
}
|
||||
__classPrivateFieldGet(this, _MessageStream_rejectConnectedPromise, "f").call(this, error);
|
||||
__classPrivateFieldGet(this, _MessageStream_rejectEndPromise, "f").call(this, error);
|
||||
this._emit('end');
|
||||
return;
|
||||
}
|
||||
if (event === 'error') {
|
||||
// NOTE: _emit('error', error) should only be called from #handleError().
|
||||
const error = args[0];
|
||||
if (!__classPrivateFieldGet(this, _MessageStream_catchingPromiseCreated, "f") && !listeners?.length) {
|
||||
// Trigger an unhandled rejection if the user hasn't registered any error handlers.
|
||||
// If you are seeing stack traces here, make sure to handle errors via either:
|
||||
// - runner.on('error', () => ...)
|
||||
// - await runner.done()
|
||||
// - await runner.final...()
|
||||
// - etc.
|
||||
Promise.reject(error);
|
||||
}
|
||||
__classPrivateFieldGet(this, _MessageStream_rejectConnectedPromise, "f").call(this, error);
|
||||
__classPrivateFieldGet(this, _MessageStream_rejectEndPromise, "f").call(this, error);
|
||||
this._emit('end');
|
||||
}
|
||||
}
|
||||
_emitFinal() {
|
||||
const finalMessage = this.receivedMessages.at(-1);
|
||||
if (finalMessage) {
|
||||
this._emit('finalMessage', __classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_getFinalMessage).call(this));
|
||||
}
|
||||
}
|
||||
async _fromReadableStream(readableStream, options) {
|
||||
const signal = options?.signal;
|
||||
if (signal) {
|
||||
if (signal.aborted)
|
||||
this.controller.abort();
|
||||
signal.addEventListener('abort', () => this.controller.abort());
|
||||
}
|
||||
__classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_beginRequest).call(this);
|
||||
this._connected(null);
|
||||
const stream = Stream.fromReadableStream(readableStream, this.controller);
|
||||
for await (const event of stream) {
|
||||
__classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_addStreamEvent).call(this, event);
|
||||
}
|
||||
if (stream.controller.signal?.aborted) {
|
||||
throw new APIUserAbortError();
|
||||
}
|
||||
__classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_endRequest).call(this);
|
||||
}
|
||||
[(_MessageStream_currentMessageSnapshot = new WeakMap(), _MessageStream_connectedPromise = new WeakMap(), _MessageStream_resolveConnectedPromise = new WeakMap(), _MessageStream_rejectConnectedPromise = new WeakMap(), _MessageStream_endPromise = new WeakMap(), _MessageStream_resolveEndPromise = new WeakMap(), _MessageStream_rejectEndPromise = new WeakMap(), _MessageStream_listeners = new WeakMap(), _MessageStream_ended = new WeakMap(), _MessageStream_errored = new WeakMap(), _MessageStream_aborted = new WeakMap(), _MessageStream_catchingPromiseCreated = new WeakMap(), _MessageStream_response = new WeakMap(), _MessageStream_request_id = new WeakMap(), _MessageStream_handleError = new WeakMap(), _MessageStream_instances = new WeakSet(), _MessageStream_getFinalMessage = function _MessageStream_getFinalMessage() {
|
||||
if (this.receivedMessages.length === 0) {
|
||||
throw new AnthropicError('stream ended without producing a Message with role=assistant');
|
||||
}
|
||||
return this.receivedMessages.at(-1);
|
||||
}, _MessageStream_getFinalText = function _MessageStream_getFinalText() {
|
||||
if (this.receivedMessages.length === 0) {
|
||||
throw new AnthropicError('stream ended without producing a Message with role=assistant');
|
||||
}
|
||||
const textBlocks = this.receivedMessages
|
||||
.at(-1)
|
||||
.content.filter((block) => block.type === 'text')
|
||||
.map((block) => block.text);
|
||||
if (textBlocks.length === 0) {
|
||||
throw new AnthropicError('stream ended without producing a content block with type=text');
|
||||
}
|
||||
return textBlocks.join(' ');
|
||||
}, _MessageStream_beginRequest = function _MessageStream_beginRequest() {
|
||||
if (this.ended)
|
||||
return;
|
||||
__classPrivateFieldSet(this, _MessageStream_currentMessageSnapshot, undefined, "f");
|
||||
}, _MessageStream_addStreamEvent = function _MessageStream_addStreamEvent(event) {
|
||||
if (this.ended)
|
||||
return;
|
||||
const messageSnapshot = __classPrivateFieldGet(this, _MessageStream_instances, "m", _MessageStream_accumulateMessage).call(this, event);
|
||||
this._emit('streamEvent', event, messageSnapshot);
|
||||
switch (event.type) {
|
||||
case 'content_block_delta': {
|
||||
const content = messageSnapshot.content.at(-1);
|
||||
switch (event.delta.type) {
|
||||
case 'text_delta': {
|
||||
if (content.type === 'text') {
|
||||
this._emit('text', event.delta.text, content.text || '');
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'citations_delta': {
|
||||
if (content.type === 'text') {
|
||||
this._emit('citation', event.delta.citation, content.citations ?? []);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'input_json_delta': {
|
||||
if (content.type === 'tool_use' && content.input) {
|
||||
this._emit('inputJson', event.delta.partial_json, content.input);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'thinking_delta': {
|
||||
if (content.type === 'thinking') {
|
||||
this._emit('thinking', event.delta.thinking, content.thinking);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'signature_delta': {
|
||||
// we don't emit anything special in this case.
|
||||
break;
|
||||
}
|
||||
default:
|
||||
checkNever(event.delta);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'message_stop': {
|
||||
this._addMessageParam(messageSnapshot);
|
||||
this._addMessage(messageSnapshot, true);
|
||||
break;
|
||||
}
|
||||
case 'content_block_stop': {
|
||||
this._emit('contentBlock', messageSnapshot.content.at(-1));
|
||||
break;
|
||||
}
|
||||
case 'message_start': {
|
||||
__classPrivateFieldSet(this, _MessageStream_currentMessageSnapshot, messageSnapshot, "f");
|
||||
break;
|
||||
}
|
||||
case 'content_block_start':
|
||||
case 'message_delta':
|
||||
break;
|
||||
}
|
||||
}, _MessageStream_endRequest = function _MessageStream_endRequest() {
|
||||
if (this.ended) {
|
||||
throw new AnthropicError(`stream has ended, this shouldn't happen`);
|
||||
}
|
||||
const snapshot = __classPrivateFieldGet(this, _MessageStream_currentMessageSnapshot, "f");
|
||||
if (!snapshot) {
|
||||
throw new AnthropicError(`request ended without sending any chunks`);
|
||||
}
|
||||
__classPrivateFieldSet(this, _MessageStream_currentMessageSnapshot, undefined, "f");
|
||||
return snapshot;
|
||||
}, _MessageStream_accumulateMessage = function _MessageStream_accumulateMessage(event) {
|
||||
let snapshot = __classPrivateFieldGet(this, _MessageStream_currentMessageSnapshot, "f");
|
||||
if (event.type === 'message_start') {
|
||||
if (snapshot) {
|
||||
throw new AnthropicError(`Unexpected event order, got ${event.type} before receiving "message_stop"`);
|
||||
}
|
||||
return event.message;
|
||||
}
|
||||
if (!snapshot) {
|
||||
throw new AnthropicError(`Unexpected event order, got ${event.type} before "message_start"`);
|
||||
}
|
||||
switch (event.type) {
|
||||
case 'message_stop':
|
||||
return snapshot;
|
||||
case 'message_delta':
|
||||
snapshot.stop_reason = event.delta.stop_reason;
|
||||
snapshot.stop_sequence = event.delta.stop_sequence;
|
||||
snapshot.usage.output_tokens = event.usage.output_tokens;
|
||||
return snapshot;
|
||||
case 'content_block_start':
|
||||
snapshot.content.push(event.content_block);
|
||||
return snapshot;
|
||||
case 'content_block_delta': {
|
||||
const snapshotContent = snapshot.content.at(event.index);
|
||||
switch (event.delta.type) {
|
||||
case 'text_delta': {
|
||||
if (snapshotContent?.type === 'text') {
|
||||
snapshotContent.text += event.delta.text;
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'citations_delta': {
|
||||
if (snapshotContent?.type === 'text') {
|
||||
snapshotContent.citations ?? (snapshotContent.citations = []);
|
||||
snapshotContent.citations.push(event.delta.citation);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'input_json_delta': {
|
||||
if (snapshotContent?.type === 'tool_use') {
|
||||
// we need to keep track of the raw JSON string as well so that we can
|
||||
// re-parse it for each delta, for now we just store it as an untyped
|
||||
// non-enumerable property on the snapshot
|
||||
let jsonBuf = snapshotContent[JSON_BUF_PROPERTY] || '';
|
||||
jsonBuf += event.delta.partial_json;
|
||||
Object.defineProperty(snapshotContent, JSON_BUF_PROPERTY, {
|
||||
value: jsonBuf,
|
||||
enumerable: false,
|
||||
writable: true,
|
||||
});
|
||||
if (jsonBuf) {
|
||||
snapshotContent.input = partialParse(jsonBuf);
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'thinking_delta': {
|
||||
if (snapshotContent?.type === 'thinking') {
|
||||
snapshotContent.thinking += event.delta.thinking;
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'signature_delta': {
|
||||
if (snapshotContent?.type === 'thinking') {
|
||||
snapshotContent.signature += event.delta.signature;
|
||||
}
|
||||
break;
|
||||
}
|
||||
default:
|
||||
checkNever(event.delta);
|
||||
}
|
||||
return snapshot;
|
||||
}
|
||||
case 'content_block_stop':
|
||||
return snapshot;
|
||||
}
|
||||
}, Symbol.asyncIterator)]() {
|
||||
const pushQueue = [];
|
||||
const readQueue = [];
|
||||
let done = false;
|
||||
this.on('streamEvent', (event) => {
|
||||
const reader = readQueue.shift();
|
||||
if (reader) {
|
||||
reader.resolve(event);
|
||||
}
|
||||
else {
|
||||
pushQueue.push(event);
|
||||
}
|
||||
});
|
||||
this.on('end', () => {
|
||||
done = true;
|
||||
for (const reader of readQueue) {
|
||||
reader.resolve(undefined);
|
||||
}
|
||||
readQueue.length = 0;
|
||||
});
|
||||
this.on('abort', (err) => {
|
||||
done = true;
|
||||
for (const reader of readQueue) {
|
||||
reader.reject(err);
|
||||
}
|
||||
readQueue.length = 0;
|
||||
});
|
||||
this.on('error', (err) => {
|
||||
done = true;
|
||||
for (const reader of readQueue) {
|
||||
reader.reject(err);
|
||||
}
|
||||
readQueue.length = 0;
|
||||
});
|
||||
return {
|
||||
next: async () => {
|
||||
if (!pushQueue.length) {
|
||||
if (done) {
|
||||
return { value: undefined, done: true };
|
||||
}
|
||||
return new Promise((resolve, reject) => readQueue.push({ resolve, reject })).then((chunk) => (chunk ? { value: chunk, done: false } : { value: undefined, done: true }));
|
||||
}
|
||||
const chunk = pushQueue.shift();
|
||||
return { value: chunk, done: false };
|
||||
},
|
||||
return: async () => {
|
||||
this.abort();
|
||||
return { value: undefined, done: true };
|
||||
},
|
||||
};
|
||||
}
|
||||
toReadableStream() {
|
||||
const stream = new Stream(this[Symbol.asyncIterator].bind(this), this.controller);
|
||||
return stream.toReadableStream();
|
||||
}
|
||||
}
|
||||
// used to ensure exhaustive case matching without throwing a runtime error
|
||||
function checkNever(x) { }
|
||||
//# sourceMappingURL=MessageStream.mjs.map
|
||||
1
vendor/sdk/lib/MessageStream.mjs.map
vendored
Normal file
1
vendor/sdk/lib/MessageStream.mjs.map
vendored
Normal file
File diff suppressed because one or more lines are too long
Reference in New Issue
Block a user