mirror of
https://github.com/RooVetGit/Roo-Code.git
synced 2026-09-07 08:26:51 +00:00
fix: handle UTF-8 boundary issues in streaming responses
- Add UTF8StreamDecoder utility to properly handle multi-byte UTF-8 characters split across chunk boundaries - Integrate decoder into OpenAI and BaseOpenAiCompatibleProvider to fix garbled output with large files - Add comprehensive tests for UTF-8 boundary handling - Fixes issue #8787 where vLLM outputs showed garbled characters with large file outputs
This commit is contained in:
parent
8187a8e189
commit
92dd39b77c
4 changed files with 433 additions and 11 deletions
|
|
@ -6,6 +6,7 @@ import type { ModelInfo } from "@roo-code/types"
|
|||
import type { ApiHandlerOptions } from "../../shared/api"
|
||||
import { ApiStream } from "../transform/stream"
|
||||
import { convertToOpenAiMessages } from "../transform/openai-format"
|
||||
import { UTF8StreamDecoder } from "../utils/utf8-stream-decoder"
|
||||
|
||||
import type { SingleCompletionHandler, ApiHandlerCreateMessageMetadata } from "../index"
|
||||
import { DEFAULT_HEADERS } from "./constants"
|
||||
|
|
@ -99,13 +100,20 @@ export abstract class BaseOpenAiCompatibleProvider<ModelName extends string>
|
|||
): ApiStream {
|
||||
const stream = await this.createStream(systemPrompt, messages, metadata)
|
||||
|
||||
// Create UTF-8 decoder for handling large outputs properly
|
||||
const utf8Decoder = new UTF8StreamDecoder()
|
||||
|
||||
for await (const chunk of stream) {
|
||||
const delta = chunk.choices[0]?.delta
|
||||
|
||||
if (delta?.content) {
|
||||
yield {
|
||||
type: "text",
|
||||
text: delta.content,
|
||||
// Decode the content properly to handle UTF-8 boundary issues
|
||||
const decodedContent = utf8Decoder.decode(delta.content)
|
||||
if (decodedContent) {
|
||||
yield {
|
||||
type: "text",
|
||||
text: decodedContent,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -117,6 +125,15 @@ export abstract class BaseOpenAiCompatibleProvider<ModelName extends string>
|
|||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Finalize any remaining buffered content
|
||||
const finalContent = utf8Decoder.finalize()
|
||||
if (finalContent) {
|
||||
yield {
|
||||
type: "text",
|
||||
text: finalContent,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async completePrompt(prompt: string): Promise<string> {
|
||||
|
|
|
|||
|
|
@ -13,6 +13,7 @@ import {
|
|||
import type { ApiHandlerOptions } from "../../shared/api"
|
||||
|
||||
import { XmlMatcher } from "../../utils/xml-matcher"
|
||||
import { UTF8StreamDecoder } from "../utils/utf8-stream-decoder"
|
||||
|
||||
import { convertToOpenAiMessages } from "../transform/openai-format"
|
||||
import { convertToR1Format } from "../transform/r1-format"
|
||||
|
|
@ -188,21 +189,32 @@ export class OpenAiHandler extends BaseProvider implements SingleCompletionHandl
|
|||
}) as const,
|
||||
)
|
||||
|
||||
// Create UTF-8 decoder for handling large outputs properly
|
||||
const utf8Decoder = new UTF8StreamDecoder()
|
||||
|
||||
let lastUsage
|
||||
|
||||
for await (const chunk of stream) {
|
||||
const delta = chunk.choices[0]?.delta ?? {}
|
||||
|
||||
if (delta.content) {
|
||||
for (const chunk of matcher.update(delta.content)) {
|
||||
yield chunk
|
||||
// Decode the content properly to handle UTF-8 boundary issues
|
||||
const decodedContent = utf8Decoder.decode(delta.content)
|
||||
if (decodedContent) {
|
||||
for (const chunk of matcher.update(decodedContent)) {
|
||||
yield chunk
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if ("reasoning_content" in delta && delta.reasoning_content) {
|
||||
yield {
|
||||
type: "reasoning",
|
||||
text: (delta.reasoning_content as string | undefined) || "",
|
||||
// Also decode reasoning content properly
|
||||
const decodedReasoning = utf8Decoder.decode(delta.reasoning_content as string)
|
||||
if (decodedReasoning) {
|
||||
yield {
|
||||
type: "reasoning",
|
||||
text: decodedReasoning,
|
||||
}
|
||||
}
|
||||
}
|
||||
if (chunk.usage) {
|
||||
|
|
@ -210,6 +222,14 @@ export class OpenAiHandler extends BaseProvider implements SingleCompletionHandl
|
|||
}
|
||||
}
|
||||
|
||||
// Finalize any remaining buffered content
|
||||
const finalContent = utf8Decoder.finalize()
|
||||
if (finalContent) {
|
||||
for (const chunk of matcher.update(finalContent)) {
|
||||
yield chunk
|
||||
}
|
||||
}
|
||||
|
||||
for (const chunk of matcher.final()) {
|
||||
yield chunk
|
||||
}
|
||||
|
|
@ -386,12 +406,19 @@ export class OpenAiHandler extends BaseProvider implements SingleCompletionHandl
|
|||
}
|
||||
|
||||
private async *handleStreamResponse(stream: AsyncIterable<OpenAI.Chat.Completions.ChatCompletionChunk>): ApiStream {
|
||||
// Create UTF-8 decoder for handling large outputs properly
|
||||
const utf8Decoder = new UTF8StreamDecoder()
|
||||
|
||||
for await (const chunk of stream) {
|
||||
const delta = chunk.choices[0]?.delta
|
||||
if (delta?.content) {
|
||||
yield {
|
||||
type: "text",
|
||||
text: delta.content,
|
||||
// Decode the content properly to handle UTF-8 boundary issues
|
||||
const decodedContent = utf8Decoder.decode(delta.content)
|
||||
if (decodedContent) {
|
||||
yield {
|
||||
type: "text",
|
||||
text: decodedContent,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -403,6 +430,15 @@ export class OpenAiHandler extends BaseProvider implements SingleCompletionHandl
|
|||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Finalize any remaining buffered content
|
||||
const finalContent = utf8Decoder.finalize()
|
||||
if (finalContent) {
|
||||
yield {
|
||||
type: "text",
|
||||
text: finalContent,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private _getUrlHost(baseUrl?: string): string {
|
||||
|
|
|
|||
197
src/api/utils/__tests__/utf8-stream-decoder.test.ts
Normal file
197
src/api/utils/__tests__/utf8-stream-decoder.test.ts
Normal file
|
|
@ -0,0 +1,197 @@
|
|||
import { describe, it, expect, beforeEach } from "vitest"
|
||||
import { UTF8StreamDecoder } from "../utf8-stream-decoder"
|
||||
|
||||
describe("UTF8StreamDecoder", () => {
|
||||
let decoder: UTF8StreamDecoder
|
||||
|
||||
beforeEach(() => {
|
||||
decoder = new UTF8StreamDecoder()
|
||||
})
|
||||
|
||||
describe("decode", () => {
|
||||
it("should handle complete ASCII strings", () => {
|
||||
const result = decoder.decode("Hello World")
|
||||
expect(result).toBe("Hello World")
|
||||
})
|
||||
|
||||
it("should handle complete UTF-8 strings", () => {
|
||||
const result = decoder.decode("Hello 世界 🌍")
|
||||
expect(result).toBe("Hello 世界 🌍")
|
||||
})
|
||||
|
||||
it("should handle multi-byte UTF-8 characters split across chunks", () => {
|
||||
// "世" (U+4E16) in UTF-8 is 0xE4 0xB8 0x96
|
||||
// Split it across two chunks
|
||||
const chunk1 = new Uint8Array([0x48, 0x65, 0x6c, 0x6c, 0x6f, 0x20, 0xe4]) // "Hello " + first byte of "世"
|
||||
const chunk2 = new Uint8Array([0xb8, 0x96]) // remaining bytes of "世"
|
||||
|
||||
const result1 = decoder.decode(chunk1)
|
||||
expect(result1).toBe("Hello ") // Should only decode complete characters
|
||||
|
||||
const result2 = decoder.decode(chunk2)
|
||||
expect(result2).toBe("世") // Should complete the character
|
||||
})
|
||||
|
||||
it("should handle 4-byte emoji split across chunks", () => {
|
||||
// "🌍" (U+1F30D) in UTF-8 is 0xF0 0x9F 0x8C 0x8D
|
||||
// Split it across multiple chunks
|
||||
const chunk1 = new Uint8Array([0x48, 0x69, 0x20, 0xf0]) // "Hi " + first byte
|
||||
const chunk2 = new Uint8Array([0x9f, 0x8c]) // middle bytes
|
||||
const chunk3 = new Uint8Array([0x8d, 0x21]) // last byte + "!"
|
||||
|
||||
const result1 = decoder.decode(chunk1)
|
||||
expect(result1).toBe("Hi ") // Should only decode complete characters
|
||||
|
||||
const result2 = decoder.decode(chunk2)
|
||||
expect(result2).toBe("") // Still incomplete
|
||||
|
||||
const result3 = decoder.decode(chunk3)
|
||||
expect(result3).toBe("🌍!") // Should complete the emoji and include the exclamation
|
||||
})
|
||||
|
||||
it("should handle string chunks with potential partial sequences", () => {
|
||||
// Simulate a string that ends with a partial UTF-8 sequence marker
|
||||
const chunk1 = "Hello 世"
|
||||
const chunk2 = "界 World"
|
||||
|
||||
const result1 = decoder.decode(chunk1)
|
||||
const result2 = decoder.decode(chunk2)
|
||||
|
||||
expect(result1 + result2).toBe("Hello 世界 World")
|
||||
})
|
||||
|
||||
it("should handle replacement characters properly", () => {
|
||||
// Test with actual replacement characters (U+FFFD)
|
||||
const chunk = "Hello \uFFFD World"
|
||||
const result = decoder.decode(chunk)
|
||||
expect(result).toBe("Hello \uFFFD World")
|
||||
})
|
||||
|
||||
it("should handle replacement characters in the middle of text", () => {
|
||||
// Replacement characters in the middle should be preserved
|
||||
const chunk = "Hello \uFFFD World"
|
||||
const result = decoder.decode(chunk)
|
||||
expect(result).toBe("Hello \uFFFD World")
|
||||
})
|
||||
|
||||
it("should handle multiple replacement characters", () => {
|
||||
// Multiple replacement characters might indicate encoding issues
|
||||
// but we should preserve them as they might be intentional
|
||||
const chunk1 = "Hello World"
|
||||
const chunk2 = " Test"
|
||||
|
||||
const result1 = decoder.decode(chunk1)
|
||||
expect(result1).toBe("Hello World")
|
||||
|
||||
const result2 = decoder.decode(chunk2)
|
||||
expect(result2).toBe(" Test")
|
||||
})
|
||||
|
||||
it("should handle empty chunks", () => {
|
||||
const result = decoder.decode("")
|
||||
expect(result).toBe("")
|
||||
})
|
||||
|
||||
it("should handle Uint8Array empty chunks", () => {
|
||||
const result = decoder.decode(new Uint8Array(0))
|
||||
expect(result).toBe("")
|
||||
})
|
||||
})
|
||||
|
||||
describe("finalize", () => {
|
||||
it("should return empty string when no buffered content", () => {
|
||||
const result = decoder.finalize()
|
||||
expect(result).toBe("")
|
||||
})
|
||||
|
||||
it("should decode remaining buffered content", () => {
|
||||
// Send an incomplete sequence
|
||||
const chunk = new Uint8Array([0x48, 0x65, 0x6c, 0x6c, 0x6f, 0x20, 0xe4, 0xb8]) // "Hello " + partial "世"
|
||||
decoder.decode(chunk)
|
||||
|
||||
// Finalize should attempt to decode what's left (may produce replacement character)
|
||||
const result = decoder.finalize()
|
||||
expect(result.length).toBeGreaterThan(0) // Should produce something (likely with replacement char)
|
||||
})
|
||||
|
||||
it("should clear buffer after finalize", () => {
|
||||
const chunk = new Uint8Array([0xe4]) // Partial character
|
||||
decoder.decode(chunk)
|
||||
|
||||
decoder.finalize()
|
||||
const secondFinalize = decoder.finalize()
|
||||
expect(secondFinalize).toBe("") // Buffer should be empty
|
||||
})
|
||||
})
|
||||
|
||||
describe("reset", () => {
|
||||
it("should clear the buffer", () => {
|
||||
// Add some partial data
|
||||
const chunk = new Uint8Array([0xe4, 0xb8]) // Partial character
|
||||
decoder.decode(chunk)
|
||||
|
||||
// Reset
|
||||
decoder.reset()
|
||||
|
||||
// Should start fresh
|
||||
const result = decoder.decode("Hello")
|
||||
expect(result).toBe("Hello")
|
||||
|
||||
// Finalize should return nothing
|
||||
const final = decoder.finalize()
|
||||
expect(final).toBe("")
|
||||
})
|
||||
})
|
||||
|
||||
describe("large file handling", () => {
|
||||
it("should handle large text with many UTF-8 characters", () => {
|
||||
// Simulate a large file with mixed content
|
||||
const largeText = "初めまして、私は人工知能です。" + "世界は美しい。".repeat(100) + "🌍🌎🌏"
|
||||
|
||||
// Split into random chunks to simulate streaming
|
||||
const chunkSize = 17 // Prime number to ensure we split across character boundaries
|
||||
const chunks: string[] = []
|
||||
for (let i = 0; i < largeText.length; i += chunkSize) {
|
||||
chunks.push(largeText.slice(i, i + chunkSize))
|
||||
}
|
||||
|
||||
// Decode all chunks
|
||||
let result = ""
|
||||
for (const chunk of chunks) {
|
||||
result += decoder.decode(chunk)
|
||||
}
|
||||
result += decoder.finalize()
|
||||
|
||||
// Should reconstruct the original text
|
||||
expect(result).toBe(largeText)
|
||||
})
|
||||
|
||||
it("should handle simulated vLLM output with potential garbling", () => {
|
||||
// Simulate what might come from vLLM with large outputs
|
||||
const chunks = [
|
||||
"def process_data(items):\n",
|
||||
' """Process a list of items',
|
||||
" with special handling for UTF-8",
|
||||
" characters like 你好", // Chinese characters might be split
|
||||
'世界"""\n result = []\n',
|
||||
" for item in items:\n",
|
||||
" # Handle special chars: €£¥",
|
||||
"🔧🔨\n",
|
||||
" result.append(transform(item))\n",
|
||||
" return result",
|
||||
]
|
||||
|
||||
let decoded = ""
|
||||
for (const chunk of chunks) {
|
||||
decoded += decoder.decode(chunk)
|
||||
}
|
||||
decoded += decoder.finalize()
|
||||
|
||||
// Should contain all the expected content without garbling
|
||||
expect(decoded).toContain("你好世界")
|
||||
expect(decoded).toContain("€£¥")
|
||||
expect(decoded).toContain("🔧🔨")
|
||||
expect(decoded).not.toContain("\uFFFD") // Should not have replacement characters
|
||||
})
|
||||
})
|
||||
})
|
||||
172
src/api/utils/utf8-stream-decoder.ts
Normal file
172
src/api/utils/utf8-stream-decoder.ts
Normal file
|
|
@ -0,0 +1,172 @@
|
|||
/**
|
||||
* UTF-8 Stream Decoder
|
||||
*
|
||||
* This utility handles proper UTF-8 decoding for streaming responses,
|
||||
* ensuring that multi-byte UTF-8 characters split across chunk boundaries
|
||||
* are properly handled without producing garbled output.
|
||||
*
|
||||
* This fixes issues with large outputs from models like vLLM where
|
||||
* characters can be split across streaming chunks.
|
||||
*/
|
||||
export class UTF8StreamDecoder {
|
||||
private buffer: Uint8Array = new Uint8Array(0)
|
||||
private decoder = new TextDecoder("utf-8", { fatal: false, ignoreBOM: true })
|
||||
|
||||
/**
|
||||
* Decodes a chunk of data, handling partial UTF-8 sequences
|
||||
* @param chunk - The chunk to decode (string or Uint8Array)
|
||||
* @returns The decoded string, with partial sequences buffered
|
||||
*/
|
||||
decode(chunk: string | Uint8Array): string {
|
||||
// If chunk is already a string, check if it needs special handling
|
||||
if (typeof chunk === "string") {
|
||||
// Check for potential UTF-8 issues in the string
|
||||
// This can happen when the OpenAI SDK doesn't properly handle chunk boundaries
|
||||
return this.handleStringChunk(chunk)
|
||||
}
|
||||
|
||||
// Convert to Uint8Array if needed
|
||||
const bytes = chunk instanceof Uint8Array ? chunk : new TextEncoder().encode(chunk)
|
||||
|
||||
// Combine with any buffered bytes from previous chunks
|
||||
const combined = this.combineBuffers(this.buffer, bytes)
|
||||
|
||||
// Find the last complete UTF-8 character boundary
|
||||
const lastCompleteIndex = this.findLastCompleteCharBoundary(combined)
|
||||
|
||||
if (lastCompleteIndex === -1) {
|
||||
// No complete characters, buffer everything
|
||||
this.buffer = combined
|
||||
return ""
|
||||
}
|
||||
|
||||
// Decode the complete portion
|
||||
const complete = combined.slice(0, lastCompleteIndex)
|
||||
const decoded = this.decoder.decode(complete, { stream: true })
|
||||
|
||||
// Buffer the incomplete portion
|
||||
this.buffer = combined.slice(lastCompleteIndex)
|
||||
|
||||
return decoded
|
||||
}
|
||||
|
||||
/**
|
||||
* Handles string chunks that might contain garbled characters
|
||||
* @param chunk - The string chunk to process
|
||||
* @returns The cleaned string
|
||||
*/
|
||||
private handleStringChunk(chunk: string): string {
|
||||
// For string chunks, we mainly need to handle cases where
|
||||
// the chunk ends with an incomplete UTF-8 sequence
|
||||
// We'll be conservative and only buffer incomplete sequences at the end
|
||||
|
||||
// If the chunk appears to end with a partial UTF-8 sequence
|
||||
// (indicated by specific byte patterns when re-encoded)
|
||||
try {
|
||||
const encoder = new TextEncoder()
|
||||
const bytes = encoder.encode(chunk)
|
||||
|
||||
// Check if the last few bytes could be a partial UTF-8 sequence
|
||||
const lastValid = this.findLastCompleteCharBoundary(bytes)
|
||||
if (lastValid > 0 && lastValid < bytes.length) {
|
||||
// We have a partial sequence at the end
|
||||
// Decode only the complete portion
|
||||
const validBytes = bytes.slice(0, lastValid)
|
||||
const decoded = this.decoder.decode(validBytes, { stream: true })
|
||||
|
||||
// Buffer the incomplete portion for the next chunk
|
||||
this.buffer = bytes.slice(lastValid)
|
||||
return decoded
|
||||
}
|
||||
} catch (e) {
|
||||
// If re-encoding fails, just return the chunk as-is
|
||||
// The OpenAI SDK should already be handling most encoding properly
|
||||
}
|
||||
|
||||
return chunk
|
||||
}
|
||||
|
||||
/**
|
||||
* Finalizes decoding, processing any remaining buffered bytes
|
||||
* @returns Any remaining decoded text
|
||||
*/
|
||||
finalize(): string {
|
||||
if (this.buffer.length === 0) {
|
||||
return ""
|
||||
}
|
||||
|
||||
// Decode whatever is left in the buffer
|
||||
// Use 'replace' mode to handle any incomplete sequences
|
||||
const decoded = new TextDecoder("utf-8", { fatal: false }).decode(this.buffer)
|
||||
this.buffer = new Uint8Array(0)
|
||||
return decoded
|
||||
}
|
||||
|
||||
/**
|
||||
* Combines two Uint8Arrays
|
||||
*/
|
||||
private combineBuffers(a: Uint8Array, b: Uint8Array): Uint8Array {
|
||||
const combined = new Uint8Array(a.length + b.length)
|
||||
combined.set(a, 0)
|
||||
combined.set(b, a.length)
|
||||
return combined
|
||||
}
|
||||
|
||||
/**
|
||||
* Finds the last complete UTF-8 character boundary in a byte array
|
||||
* @param bytes - The byte array to check
|
||||
* @returns The index after the last complete character, or -1 if none
|
||||
*/
|
||||
private findLastCompleteCharBoundary(bytes: Uint8Array): number {
|
||||
if (bytes.length === 0) return -1
|
||||
|
||||
// Work backwards to find the last complete UTF-8 character
|
||||
for (let i = bytes.length; i > 0; i--) {
|
||||
const byte = bytes[i - 1]
|
||||
|
||||
// Check if this is a valid UTF-8 sequence start or single byte
|
||||
if ((byte & 0x80) === 0) {
|
||||
// Single byte character (0xxxxxxx)
|
||||
return i
|
||||
} else if ((byte & 0xe0) === 0xc0) {
|
||||
// 2-byte sequence start (110xxxxx)
|
||||
if (i + 1 <= bytes.length && this.isValidContinuation(bytes, i, 1)) {
|
||||
return i + 1
|
||||
}
|
||||
} else if ((byte & 0xf0) === 0xe0) {
|
||||
// 3-byte sequence start (1110xxxx)
|
||||
if (i + 2 <= bytes.length && this.isValidContinuation(bytes, i, 2)) {
|
||||
return i + 2
|
||||
}
|
||||
} else if ((byte & 0xf8) === 0xf0) {
|
||||
// 4-byte sequence start (11110xxx)
|
||||
if (i + 3 <= bytes.length && this.isValidContinuation(bytes, i, 3)) {
|
||||
return i + 3
|
||||
}
|
||||
}
|
||||
// If we hit a continuation byte (10xxxxxx), keep going back
|
||||
}
|
||||
|
||||
return -1
|
||||
}
|
||||
|
||||
/**
|
||||
* Checks if the continuation bytes are valid
|
||||
*/
|
||||
private isValidContinuation(bytes: Uint8Array, start: number, count: number): boolean {
|
||||
for (let i = 0; i < count; i++) {
|
||||
if (start + i >= bytes.length) return false
|
||||
const byte = bytes[start + i]
|
||||
// Continuation bytes must match 10xxxxxx
|
||||
if ((byte & 0xc0) !== 0x80) return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
/**
|
||||
* Resets the decoder state
|
||||
*/
|
||||
reset(): void {
|
||||
this.buffer = new Uint8Array(0)
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue