fix: chat stream data too long
This commit is contained in:
@@ -3,30 +3,33 @@ import qs from 'query-string';
|
||||
|
||||
const extractStreamRegx = /(data|error):\s*({.*?})(?=\n|$)/g;
|
||||
|
||||
const extractJSON = (dataStr: string) => {
|
||||
let match;
|
||||
const extractJSON = (
|
||||
dataStr: string
|
||||
): { results: any[]; remaining: string } => {
|
||||
const results: any[] = [];
|
||||
if (!dataStr) return { results, remaining: '' };
|
||||
|
||||
if (!dataStr) {
|
||||
return results;
|
||||
}
|
||||
while ((match = extractStreamRegx.exec(dataStr)) !== null) {
|
||||
try {
|
||||
const type = match[1]; // "data" or "error"
|
||||
console.log('type========', type);
|
||||
const jsonData = JSON.parse(match[2]);
|
||||
if (type === 'error') {
|
||||
results.push({ error: jsonData });
|
||||
} else {
|
||||
results.push(jsonData);
|
||||
}
|
||||
} catch (err) {
|
||||
console.error('JSON parse error:', err, 'for match:', match[2]);
|
||||
const parts = dataStr.split(/\n\n/);
|
||||
let remaining = '';
|
||||
|
||||
for (let i = 0; i < parts.length; i++) {
|
||||
const part = parts[i].trim();
|
||||
if (!part) continue;
|
||||
if (!part.startsWith('data:') && !part.startsWith('error:')) {
|
||||
remaining += part;
|
||||
continue;
|
||||
}
|
||||
|
||||
const jsonStr = part.slice(5).trim();
|
||||
try {
|
||||
const parsed = JSON.parse(jsonStr);
|
||||
results.push(parsed);
|
||||
} catch {
|
||||
remaining += part;
|
||||
}
|
||||
}
|
||||
|
||||
return results;
|
||||
return { results, remaining };
|
||||
};
|
||||
|
||||
const errorHandler = async (res: any) => {
|
||||
@@ -181,10 +184,11 @@ export const readStreamData = async (
|
||||
bufferManager.add({ error: jsonData });
|
||||
textBuffer = ''; // Clear buffer after processing error
|
||||
} else {
|
||||
extractJSON(textBuffer).forEach((data) => {
|
||||
const { results, remaining } = extractJSON(textBuffer);
|
||||
results.forEach((data) => {
|
||||
bufferManager.add(data);
|
||||
});
|
||||
textBuffer = ''; // Clear buffer after processing
|
||||
textBuffer = remaining; // Clear buffer after processing
|
||||
}
|
||||
|
||||
throttledCallback();
|
||||
@@ -195,10 +199,11 @@ export const readStreamData = async (
|
||||
if (done) {
|
||||
textBuffer += decoder.decode();
|
||||
if (textBuffer) {
|
||||
extractJSON(textBuffer).forEach((data) => bufferManager.add(data));
|
||||
const { results, remaining } = extractJSON(textBuffer);
|
||||
results.forEach((data) => bufferManager.add(data));
|
||||
}
|
||||
isReading = false;
|
||||
bufferManager.flush();
|
||||
isReading = false;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user