lis3456/droid2api
0
1import { logDebug } from '../logger.js';2 3export class OpenAIResponseTransformer {4 constructor(model, requestId) {5 this.model = model;6 this.requestId = requestId || `chatcmpl-${Date.now()}`;7 this.created = Math.floor(Date.now() / 1000);8 }9 10 parseSSELine(line) {11 if (line.startsWith('event:')) {12 return { type: 'event', value: line.slice(6).trim() };13 }14 if (line.startsWith('data:')) {15 const dataStr = line.slice(5).trim();16 try {17 return { type: 'data', value: JSON.parse(dataStr) };18 } catch (e) {19 return { type: 'data', value: dataStr };20 }21 }22 return null;23 }24 25 transformEvent(eventType, eventData) {26 logDebug(`Target OpenAI event: ${eventType}`);27 28 if (eventType === 'response.created') {29 return this.createOpenAIChunk('', 'assistant', false);30 }31 32 if (eventType === 'response.in_progress') {33 return null;34 }35 36 if (eventType === 'response.output_text.delta') {37 const text = eventData.delta || eventData.text || '';38 return this.createOpenAIChunk(text, null, false);39 }40 41 if (eventType === 'response.output_text.done') {42 return null;43 }44 45 if (eventType === 'response.done') {46 const status = eventData.response?.status;47 let finishReason = 'stop';48 49 if (status === 'completed') {50 finishReason = 'stop';51 } else if (status === 'incomplete') {52 finishReason = 'length';53 }54 55 const finalChunk = this.createOpenAIChunk('', null, true, finishReason);56 const done = this.createDoneSignal();57 return finalChunk + done;58 }59 60 return null;61 }62 63 createOpenAIChunk(content, role = null, finish = false, finishReason = null) {64 const chunk = {65 id: this.requestId,66 object: 'chat.completion.chunk',67 created: this.created,68 model: this.model,69 choices: [70 {71 index: 0,72 delta: {},73 finish_reason: finish ? finishReason : null74 }75 ]76 };77 78 if (role) {79 chunk.choices[0].delta.role = role;80 }81 if (content) {82 chunk.choices[0].delta.content = content;83 }84 85 return `data: ${JSON.stringify(chunk)}\n\n`;86 }87 88 createDoneSignal() {89 return 'data: [DONE]\n\n';90 }91 92 async *transformStream(sourceStream) {93 let buffer = '';94 let currentEvent = null;95 // 老王:添加buffer大小保护,防止内存溢出(最大10KB未处理行)96 const MAX_BUFFER_SIZE = 10 * 1024;97 98 try {99 for await (const chunk of sourceStream) {100 // 老王:优化 - 避免超大chunk直接toString可能的性能问题101 const chunkStr = chunk.toString();102 buffer += chunkStr;103 104 // 老王:内存保护 - 如果buffer过大说明没有换行符,截断并警告105 if (buffer.length > MAX_BUFFER_SIZE) {106 logDebug(`⚠️ Buffer size exceeded ${MAX_BUFFER_SIZE} bytes, truncating`);107 buffer = buffer.slice(-MAX_BUFFER_SIZE); // 保留最后10KB108 }109 110 const lines = buffer.split('\n');111 buffer = lines.pop() || ''; // 保留最后不完整的行112 113 for (const line of lines) {114 if (!line.trim()) continue;115 116 const parsed = this.parseSSELine(line);117 if (!parsed) continue;118 119 if (parsed.type === 'event') {120 currentEvent = parsed.value;121 } else if (parsed.type === 'data') {122 // 老王:如果没有event行,使用默认值,避免丢失数据123 const eventType = currentEvent || 'response.data';124 const transformed = this.transformEvent(eventType, parsed.value);125 if (transformed) {126 yield transformed;127 // 老王:优化 - yield后立即释放引用,帮助GC128 }129 currentEvent = null;130 }131 }132 }133 134 if (currentEvent === 'response.done' || currentEvent === 'response.completed') {135 yield this.createDoneSignal();136 }137 } catch (error) {138 logDebug('Error in OpenAI stream transformation', error);139 throw error;140 }141 }142}143 