Skip to content

Commit b9571cc

Browse files
committed
refactor(ai): SSE 读取内核收敛与 NDJSON 喂入纯函数+单测(观察 2026-09-05-2)
1 parent 14e3988 commit b9571cc

4 files changed

Lines changed: 113 additions & 15 deletions

File tree

‎app/src-tauri/src/ai_chat_stream.rs‎

Lines changed: 19 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -174,19 +174,35 @@ pub fn stream_sse_content(
174174
.set("Authorization", &format!("Bearer {}", client.config.api_key.trim()))
175175
.send_string(&payload.to_string())
176176
.map_err(map_status)?;
177-
let reader = BufReader::new(resp.into_reader());
177+
let (content, usage_json, cancelled, completed) =
178+
read_sse_lines(BufReader::new(resp.into_reader()), None, |d| emit(d));
179+
Ok(StreamOutcome { content, usage_json, cancelled, completed })
180+
}
181+
182+
/// 通用 SSE 读取内核(观察 2026-09-05-2:收敛与 stream_chat 高度同构的循环——
183+
/// 断流行 completed=false,调用方不得当成功;cancel 可选供聊天路径复用)。
184+
fn read_sse_lines(
185+
reader: impl std::io::BufRead,
186+
cancel: Option<&CancelFlag>,
187+
mut on_delta: impl FnMut(&str),
188+
) -> (String, Option<String>, bool, bool) {
178189
let mut content = String::new();
179190
let mut usage_json: Option<String> = None;
191+
let mut cancelled = false;
180192
let mut completed = false;
181193
for line in reader.lines() {
194+
if cancel.is_some_and(|c| c.is_cancelled()) {
195+
cancelled = true;
196+
break;
197+
}
182198
let line = match line {
183199
Ok(l) => l,
184200
Err(_) => break, // 传输中途断流:completed=false——调用方不得当成功
185201
};
186202
match parse_sse_line(&line) {
187203
SseEvent::Delta(d) => {
188204
content.push_str(&d);
189-
emit(&d);
205+
on_delta(&d);
190206
}
191207
SseEvent::Done => {
192208
completed = true;
@@ -199,7 +215,7 @@ pub fn stream_sse_content(
199215
}
200216
}
201217
}
202-
Ok(StreamOutcome { content, usage_json, cancelled: false, completed })
218+
(content, usage_json, cancelled, completed)
203219
}
204220

205221
/// HTTP 状态 → AiClientError(与 post_completions 同归一口径——四下一致)。

‎app/src-tauri/src/ai_note_refine.rs‎

Lines changed: 11 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -216,8 +216,6 @@ impl AiNoteRefineAdapter {
216216
dims: Option<&ResolvedDims>,
217217
mut on_section: impl FnMut(crate::ai_refine_protocol::AiRefineSection),
218218
) -> Result<AiRefineResponse, AiClientError> {
219-
use crate::ai_refine_protocol::parse_section_ndjson_line;
220-
221219
let mut system = self.prompt.build_system(&request.profile, dims);
222220
system.push_str("\n\n");
223221
system.push_str(crate::ai_refine_protocol::NDJSON_SYSTEM_SUFFIX);
@@ -250,20 +248,21 @@ impl AiNoteRefineAdapter {
250248
&self.client,
251249
payload,
252250
|delta| {
253-
pending.push_str(delta);
254-
while let Some(pos) = pending.find('\n') {
255-
let line: String = pending.drain(..=pos).collect();
256-
if let Some(sec) = parse_section_ndjson_line(&line) {
257-
sections.push(sec.clone());
258-
on_section(sec);
259-
}
251+
// 观察 2026-09-05-2:行缓冲收敛至 ndjson_feed 纯函数(可单测)
252+
let before = sections.len();
253+
crate::ndjson_feed::feed_ndjson(&mut pending, delta, &mut sections);
254+
for sec in sections[before..].iter().cloned() {
255+
on_section(sec);
260256
}
261257
},
262258
)?;
263259
// 末行无换行(SSE 结束前 flush)
264-
if let Some(sec) = parse_section_ndjson_line(&pending) {
265-
sections.push(sec.clone());
266-
on_section(sec);
260+
{
261+
let before = sections.len();
262+
crate::ndjson_feed::flush_ndjson(&mut pending, &mut sections);
263+
for sec in sections[before..].iter().cloned() {
264+
on_section(sec);
265+
}
267266
}
268267
// B2(审查):未收到 [DONE] 即断流——已累积节不可信(尾节可能丢失),
269268
// 整体视作失败走同拍非流式回退(禁止静默截断当成功)

‎app/src-tauri/src/lib.rs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -348,6 +348,8 @@ mod media_state;
348348
mod note_filter;
349349
mod note_filter_ai;
350350
mod note_filter_discourse;
351+
// 观察 2026-09-05-2:NDJSON 流式喂入缓冲纯函数(审查收口)
352+
mod ndjson_feed;
351353
// v0.12.0 M1(ADR-021):正文源多态——BodySource 判定 + OCR 精简过滤链
352354
mod note_body_source;
353355
mod note_filter_ocr;

‎app/src-tauri/src/ndjson_feed.rs‎

Lines changed: 81 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,81 @@
1+
//! NDJSON 流式喂入缓冲(观察 2026-09-05-2,v0.19.7 审查收口)。
2+
//!
3+
//! @ai-context: 精修逐节流式的行缓冲逻辑从 adapter 闭包抽为纯函数(可单测):
4+
//! chunk 可能切在行中/行间/CRLF 边界;`feed_ndjson` 只解析完整行,
5+
//! 残留留 pending;`flush_ndjson` 处理末行无换行与尾随垃圾。
6+
7+
use crate::ai_refine_protocol::{AiRefineSection, parse_section_ndjson_line};
8+
9+
/// 喂入增量文本:解析全部完整行并入 sink(未换行残留在 pending)。
10+
pub fn feed_ndjson(pending: &mut String, chunk: &str, sink: &mut Vec<AiRefineSection>) {
11+
pending.push_str(chunk);
12+
while let Some(pos) = pending.find('\n') {
13+
let line: String = pending.drain(..=pos).collect();
14+
if let Some(sec) = parse_section_ndjson_line(&line) {
15+
sink.push(sec);
16+
}
17+
}
18+
}
19+
20+
/// 流结束 flush:解析剩余(末行无换行/模型尾部垃圾行忽略)。
21+
pub fn flush_ndjson(pending: &mut String, sink: &mut Vec<AiRefineSection>) {
22+
if pending.is_empty() {
23+
return;
24+
}
25+
let tail = std::mem::take(pending);
26+
if let Some(sec) = parse_section_ndjson_line(&tail) {
27+
sink.push(sec);
28+
}
29+
}
30+
31+
#[cfg(test)]
32+
mod tests {
33+
use super::*;
34+
35+
fn sec(heading: &str) -> AiRefineSection {
36+
// 解析侧只关心 heading/blocks 形状——空 blocks 元素同构
37+
let line = format!(r#"{{"heading":"{}","blocks":[]}}"#, heading);
38+
parse_section_ndjson_line(&line).expect("构造节失败")
39+
}
40+
41+
#[test]
42+
fn chunk_split_across_lines_accumulates() {
43+
let mut pending = String::new();
44+
let mut sink = Vec::new();
45+
// 第 1 个对象被切成两半 + 第 2 个对象完整
46+
feed_ndjson(&mut pending, r#"{"heading":"节A","blocks":[]}"#, &mut sink);
47+
feed_ndjson(&mut pending, "\n", &mut sink);
48+
feed_ndjson(&mut pending, r#"{"heading":"节B","blocks":[]}"#, &mut sink);
49+
assert_eq!(sink.len(), 0, "首行未换行前不解析");
50+
feed_ndjson(&mut pending, "\n", &mut sink);
51+
assert_eq!(sink.len(), 2);
52+
assert_eq!(sink[0], sec("节A"));
53+
assert_eq!(sink[1], sec("节B"));
54+
}
55+
56+
#[test]
57+
fn crlf_and_tail_garbage_handled() {
58+
let mut pending = String::new();
59+
let mut sink = Vec::new();
60+
feed_ndjson(&mut pending, "{\"heading\":\"A\",\"blocks\":[]}\r\n{\"heading\":\"B\",\"blocks\":[]}\n", &mut sink);
61+
assert_eq!(sink.len(), 2, "CRLF 尾部 \\r 应被 trim 忽略");
62+
// 末行无换行 + 尾随解释文本 → flush 只取可解析对象
63+
feed_ndjson(&mut pending, "{\"heading\":\"C\",\"blocks\":[]}", &mut sink);
64+
feed_ndjson(&mut pending, "(以上为整理结果)", &mut sink);
65+
flush_ndjson(&mut pending, &mut sink);
66+
assert_eq!(sink.len(), 3);
67+
assert_eq!(sink[2], sec("C"));
68+
assert!(pending.is_empty());
69+
}
70+
71+
#[test]
72+
fn compact_array_lines_rejected_per_line() {
73+
// 模型输出完整数组(违反逐节约定)→ 行级全部拒绝;整包回退由
74+
// adapter 层按全文解析(见 refine_stream_ndjson)
75+
let mut pending = String::new();
76+
let mut sink = Vec::new();
77+
feed_ndjson(&mut pending, "[{\"heading\":\"A\",\"blocks\":[]}]", &mut sink);
78+
flush_ndjson(&mut pending, &mut sink);
79+
assert!(sink.is_empty());
80+
}
81+
}

0 commit comments

Comments
 (0)