import { Transform } from "stream"; import { Logger } from "./Logger"; import { ChatCompletionChunk } from "openai/resources"; import { isPunctuation } from "."; import { ChatToolCall, ChatToolCallFuntion } from "../types"; import { FunctionToolCall } from "openai/resources/beta/threads/runs/steps"; // export class ImagenResultStream extends Transform { // constructor() { // super({ readableObjectMode: true, writableObjectMode: true, objectMode: true }) // } // _transform(chunk: any, encoding: string, callback: Function) { // this.push(chunk); // callback(); // } // } export class ChatCompletionStream extends Transform { constructor() { super({ readableObjectMode: true, writableObjectMode: true, objectMode: true }) } _transform(chunk: any, encoding: string, callback: Function) { this.push(chunk); callback(); } } export class VisionResultStream extends ChatCompletionStream { constructor() { super() } } export const streamingResponseHandle = async ( response: Response, callback: (data: string | null) => void, options?: { spliteStr?: string, prefixStr?: string, endStr?: string },): Promise => { const spliteStr = options?.spliteStr || '\n\n'; const endStr = options?.endStr || '[DONE]'; const prefixStr = options?.prefixStr || 'data: '; try { if (!response.body) { Logger.error('Response body is empty'); throw new Error('Response body is empty'); } const reader = response.body.getReader(); const decoder = new TextDecoder('utf-8'); let dataBuffer = ''; while (true) { const { done, value } = await reader.read(); if (done) { callback(null) break; } dataBuffer += decoder.decode(value, { stream: true }); let dataEvents = dataBuffer.split(spliteStr); for (let i = 0; i < dataEvents.length - 1; i++) { const rawData = dataEvents[i] const lines = rawData.split('\n'); let dataLine = lines.find(line => line.startsWith(prefixStr)); if (dataLine) { dataLine = dataLine.replace(prefixStr, ''); // console.log("stream dateLine",dataLine) if (dataLine === endStr) { callback(null) } else { callback(dataLine) } } } // Keep the last incomplete event in the buffer dataBuffer = dataEvents[dataEvents.length - 1]; } return } catch (error) { return Promise.reject(error) } } export const streamingResponseCut = async ( stream: ChatCompletionStream, callback: (data: string | null) => Promise, options?: {} ) => { let contentBuf = "" for await (const iterator of stream) { const { choices } = iterator as ChatCompletionChunk const content = choices[0].delta.content || "" for (let index = 0; index < content.length; index++) { const char = content[index]; contentBuf += char if (char === '\n' || isPunctuation(char)) { await callback(contentBuf) contentBuf = "" } } } if (contentBuf) { await callback(contentBuf) } } export const streamResponseToolCallsFormat = async (stream: ChatCompletionStream): Promise => { const toolCalls: ChatToolCall[] = [] for await (const iterator of stream) { const { choices } = iterator as ChatCompletionChunk const delta = choices[0].delta delta.tool_calls?.forEach((toolCall, i) => { const name = toolCall.function?.name || "" const args = toolCall.function?.arguments ? JSON.parse(toolCall.function?.arguments || "") : {} if (!name) return let item:ChatToolCall = { id: toolCall.id || "", function: { name: toolCall.function?.name || "", arguments: args } } toolCalls.push(item) }) } return toolCalls } export const streamResponseFormat = async ( stream: ChatCompletionStream, callbacks?: { onData?: (data: any) => void, onContent?: (data: string) => void, onCut?: (data: string) => void, onToolCall?: (data: { id: string, function: { name: string, arguments?: Record } }) => void, } ): Promise => { let cutStrBuf = "" for await (const iterator of stream) { const { choices } = iterator as ChatCompletionChunk const delta = choices[0].delta if (callbacks?.onData) { callbacks.onData(choices[0]) } const content = delta.content || "" if (callbacks?.onContent && content) { callbacks.onContent(content) } for (let index = 0; index < content.length; index++) { const char = content[index]; cutStrBuf += char if (char === '\n' || char === "\t" || isPunctuation(char)) { cutStrBuf.replace("\n", "") cutStrBuf.replace("\t", "") if (callbacks?.onCut && cutStrBuf && cutStrBuf !== "\n" && cutStrBuf !== "\t") { callbacks?.onCut(cutStrBuf) } cutStrBuf = "" } } delta.tool_calls?.forEach((toolCall, i) => { callbacks?.onToolCall && callbacks.onToolCall({ id: toolCall.id || "", function: { name: toolCall.function?.name || "", arguments: toolCall.function?.arguments ? JSON.parse(toolCall.function?.arguments || "") : {} } }) }) } return } export const StreamUtils = { /** * 处理流式响应的函数。 * * 此函数的目的是从流式响应中提取所需的数据或信息。 * 它通过分析流式响应的特定字段或行为,来实现对数据的处理和解析。 * */ handle: streamingResponseHandle, /** * 分割流式响应数据。 * * 该函数用于处理流式响应数据的切割操作。通过指定切割逻辑,可以从一个持续的流式响应中提取出感兴趣的片段。 * 这对于处理大量的实时数据,比如网络传输的音频或视频流,非常有用。 */ cut: streamingResponseCut, /** * 处理工具调用的格式化函数。 * * 此函数的目的是将工具的调用结果根据特定的格式进行处理和转换, * 以便于后续的使用或展示。它接收一个特定格式的工具调用响应, * 并返回经过格式化处理的新响应。 */ toolCallsFormat: streamResponseToolCallsFormat, /** * 处理响应的格式化函数。 * * 此函数的目的是将响应结果根据特定的格式进行处理和转换, * 以便于后续的使用或展示。它接收一个特定格式的响应,并返回经过格式化处理的新响应。 */ format: streamResponseFormat }