Skip to content

Commit 0120021

Browse files
committed
update
1 parent 649a409 commit 0120021

4 files changed

Lines changed: 116 additions & 41 deletions

File tree

languages/kotlin/kotlin-codegen-reference/src/main/kotlin/ExampleClient.kt

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2147,12 +2147,20 @@ private suspend fun __handleSseRequest(
21472147
newBackoffTime = 0
21482148
val channel: ByteReadChannel = httpResponse.body()
21492149
var pendingData = ""
2150+
val decoder = Charsets.UTF_8.newDecoder()
2151+
.onMalformedInput(java.nio.charset.CodingErrorAction.REPLACE)
2152+
.onUnmappableCharacter(java.nio.charset.CodingErrorAction.REPLACE)
2153+
val byteBuffer = java.nio.ByteBuffer.allocate(bufferCapacity * 2)
21502154
while (!channel.isClosedForRead) {
2151-
val buffer = ByteBuffer.allocateDirect(bufferCapacity)
2152-
val read = channel.readAvailable(buffer)
2155+
val read = channel.readAvailable(byteBuffer)
21532156
if (read == -1) break
2154-
buffer.flip()
2155-
val input = Charsets.UTF_8.decode(buffer).toString()
2157+
byteBuffer.flip()
2158+
val charBuffer = java.nio.CharBuffer.allocate(byteBuffer.remaining())
2159+
decoder.decode(byteBuffer, charBuffer, false)
2160+
charBuffer.flip()
2161+
val input = charBuffer.toString()
2162+
byteBuffer.compact()
2163+
if (input.isEmpty()) continue
21562164
val parsedResult = __parseSseEvents("${pendingData}${input}")
21572165
pendingData = parsedResult.leftover
21582166
for (event in parsedResult.events) {

languages/kotlin/kotlin-codegen/src/_index.ts

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -629,12 +629,20 @@ private suspend fun __handleSseRequest(
629629
newBackoffTime = 0
630630
val channel: ByteReadChannel = httpResponse.body()
631631
var pendingData = ""
632+
val decoder = Charsets.UTF_8.newDecoder()
633+
.onMalformedInput(java.nio.charset.CodingErrorAction.REPLACE)
634+
.onUnmappableCharacter(java.nio.charset.CodingErrorAction.REPLACE)
635+
val byteBuffer = java.nio.ByteBuffer.allocate(bufferCapacity * 2)
632636
while (!channel.isClosedForRead) {
633-
val buffer = ByteBuffer.allocateDirect(bufferCapacity)
634-
val read = channel.readAvailable(buffer)
637+
val read = channel.readAvailable(byteBuffer)
635638
if (read == -1) break
636-
buffer.flip()
637-
val input = Charsets.UTF_8.decode(buffer).toString()
639+
byteBuffer.flip()
640+
val charBuffer = java.nio.CharBuffer.allocate(byteBuffer.remaining())
641+
decoder.decode(byteBuffer, charBuffer, false)
642+
charBuffer.flip()
643+
val input = charBuffer.toString()
644+
byteBuffer.compact()
645+
if (input.isEmpty()) continue
638646
val parsedResult = __parseSseEvents("\${pendingData}\${input}")
639647
pendingData = parsedResult.leftover
640648
for (event in parsedResult.events) {

languages/rust/rust-client/src/sse.rs

Lines changed: 45 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -239,46 +239,59 @@ impl<'a> EventSource<'a> {
239239
}
240240
self.retry_count = 0;
241241
let mut pending_data: String = "".to_string();
242+
let mut pending_bytes: Vec<u8> = Vec::new();
242243
while let Some(chunk) = ok_response.chunk().await.unwrap_or_default() {
243244
if controller.is_aborted {
244245
return SseAction::Abort;
245246
}
246-
let chunk_vec = chunk.to_vec();
247-
let data = std::str::from_utf8(chunk_vec.as_slice());
248-
match data {
249-
Ok(text) => {
250-
if !text.ends_with("\n\n") {
251-
pending_data.push_str(text);
252-
continue;
247+
pending_bytes.extend_from_slice(&chunk);
248+
let text = match std::str::from_utf8(&pending_bytes) {
249+
Ok(s) => {
250+
let s_owned = s.to_string();
251+
pending_bytes.clear();
252+
s_owned
253+
}
254+
Err(e) => {
255+
if e.error_len().is_some() {
256+
let s_owned = String::from_utf8_lossy(&pending_bytes).into_owned();
257+
pending_bytes.clear();
258+
s_owned
259+
} else {
260+
let valid_len = e.valid_up_to();
261+
let s_owned = std::str::from_utf8(&pending_bytes[..valid_len]).unwrap().to_string();
262+
pending_bytes.drain(..valid_len);
263+
s_owned
253264
}
254-
let msg_text = format!("{}{}", pending_data, text);
255-
let (messages, left_over) = sse_message_list_from_string(msg_text, false);
256-
pending_data = left_over;
257-
for message in messages {
258-
let event = message.event.unwrap_or("".to_string());
259-
match event.as_str() {
260-
"done" => {
261-
on_event(SseEvent::Close, &mut controller);
262-
return SseAction::Abort;
263-
}
264-
"message" => {
265-
on_event(
266-
SseEvent::Message(T::from_json_string(message.data)),
267-
&mut controller,
268-
);
269-
if controller.is_aborted {
270-
return SseAction::Abort;
271-
}
272-
}
273-
"" => on_event(
274-
SseEvent::Message(T::from_json_string(message.data)),
275-
&mut controller,
276-
),
277-
_ => {}
265+
}
266+
};
267+
if text.is_empty() {
268+
continue;
269+
}
270+
let msg_text = format!("{}{}", pending_data, text);
271+
let (messages, left_over) = sse_message_list_from_string(msg_text, false);
272+
pending_data = left_over;
273+
for message in messages {
274+
let event = message.event.unwrap_or("".to_string());
275+
match event.as_str() {
276+
"done" => {
277+
on_event(SseEvent::Close, &mut controller);
278+
return SseAction::Abort;
279+
}
280+
"message" => {
281+
on_event(
282+
SseEvent::Message(T::from_json_string(message.data)),
283+
&mut controller,
284+
);
285+
if controller.is_aborted {
286+
return SseAction::Abort;
278287
}
279288
}
289+
"" => on_event(
290+
SseEvent::Message(T::from_json_string(message.data)),
291+
&mut controller,
292+
),
293+
_ => {}
280294
}
281-
_ => {}
282295
}
283296
}
284297
if controller.is_aborted {

languages/swift/swift-client/Sources/ArriClient/ArriClient.swift

Lines changed: 47 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -841,13 +841,19 @@ public struct EventSource<T: ArriClientModel>: ArriCancellable {
841841
return await handleRetry()
842842
}
843843

844+
var pendingNetworkBytes = Data()
844845
do {
845846
for try await buffer in okResponse.body! {
846847
if cancelled {
847848
self.options.onClose()
848849
return
849850
}
850-
let chunk = String(buffer: buffer)
851+
let chunkData = Data(buffer: buffer)
852+
pendingNetworkBytes.append(chunkData)
853+
let chunk = extractValidUTF8(from: &pendingNetworkBytes)
854+
if chunk.isEmpty {
855+
continue
856+
}
851857
let (rawEvents, leftover) = sseEventListFromString(input: "\(pendingData)\(chunk)", debug: false)
852858
pendingData = leftover
853859
for rawEvent in rawEvents {
@@ -923,4 +929,44 @@ public struct EventSource<T: ArriClientModel>: ArriCancellable {
923929
self.retryCount += 1
924930
return await self.sendRequest()
925931
}
932+
}
933+
934+
func extractValidUTF8(from data: inout Data) -> String {
935+
guard !data.isEmpty else { return "" }
936+
var splitIndex = data.count
937+
// Check up to 4 last bytes
938+
for i in 0..<min(4, data.count) {
939+
let index = data.count - 1 - i
940+
let byte = data[index]
941+
if (byte & 0x80) == 0 {
942+
// ASCII byte, meaning any preceding sequence must be complete
943+
break
944+
}
945+
if (byte & 0xC0) == 0xC0 {
946+
// Found a leading byte
947+
let expectedLength: Int
948+
if (byte & 0xE0) == 0xC0 { expectedLength = 2 }
949+
else if (byte & 0xF0) == 0xE0 { expectedLength = 3 }
950+
else if (byte & 0xF8) == 0xF0 { expectedLength = 4 }
951+
else { expectedLength = 1 } // Invalid leading byte
952+
953+
let actualLength = i + 1
954+
if actualLength < expectedLength {
955+
splitIndex = index
956+
}
957+
break
958+
}
959+
}
960+
961+
if splitIndex == 0 {
962+
return ""
963+
}
964+
965+
let validData = Data(data.prefix(splitIndex))
966+
data.removeFirst(splitIndex)
967+
if let str = String(data: validData, encoding: .utf8) {
968+
return str
969+
}
970+
// Fallback if malformed
971+
return String(decoding: validData, as: UTF8.self)
926972
}

0 commit comments

Comments
 (0)