Skip to content

Commit b6a00a4

Browse files
committed
feat(ai): 精修片内流式——NDJSON 逐节 SSE + Delta 帧打字机(REQ-290①)
1 parent a5c912d commit b6a00a4

8 files changed

Lines changed: 353 additions & 42 deletions

File tree

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

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -143,6 +143,54 @@ pub fn extract_usage(line: &str) -> Option<String> {
143143
}
144144
}
145145

146+
/// 精修流式传输(REQ-290①):对给定 payload 置 stream:true 发 SSE,逐 delta
147+
/// 回调 emit(无取消语义——精修片幂等由重试兜底,与 chat stream_chat 区分:
148+
/// 不自动重试/不落库/无 usage 消费)。返回全部累积文本供整包回退解析。
149+
pub fn stream_sse_content(
150+
client: &AiClient,
151+
mut payload: serde_json::Value,
152+
mut emit: impl FnMut(&str),
153+
) -> Result<StreamOutcome, AiClientError> {
154+
if !client.config.is_local && client.config.api_key.trim().is_empty() {
155+
return Err(AiClientError::Auth(
156+
"未配置 API 密钥(设置页保存或配置环境变量)".to_string(),
157+
));
158+
}
159+
payload["stream"] = serde_json::json!(true);
160+
let url = chat_completions_url(&client.config.base_url);
161+
let agent = ureq::AgentBuilder::new()
162+
.timeout(std::time::Duration::from_secs(client.config.timeout_secs.max(5)))
163+
.build();
164+
let resp = agent
165+
.post(&url)
166+
.set("Content-Type", "application/json")
167+
.set("Authorization", &format!("Bearer {}", client.config.api_key.trim()))
168+
.send_string(&payload.to_string())
169+
.map_err(map_status)?;
170+
let reader = BufReader::new(resp.into_reader());
171+
let mut content = String::new();
172+
let mut usage_json: Option<String> = None;
173+
for line in reader.lines() {
174+
let line = match line {
175+
Ok(l) => l,
176+
Err(_) => break, // 传输中途断流:以已累积内容为准
177+
};
178+
match parse_sse_line(&line) {
179+
SseEvent::Delta(d) => {
180+
content.push_str(&d);
181+
emit(&d);
182+
}
183+
SseEvent::Done => break,
184+
SseEvent::Ignore => {
185+
if let Some(usage) = extract_usage(&line) {
186+
usage_json = Some(usage);
187+
}
188+
}
189+
}
190+
}
191+
Ok(StreamOutcome { content, usage_json, cancelled: false })
192+
}
193+
146194
/// HTTP 状态 → AiClientError(与 post_completions 同归一口径——四下一致)。
147195
fn map_status(e: ureq::Error) -> AiClientError {
148196
match e {

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

Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -199,6 +199,76 @@ impl AiNoteRefineAdapter {
199199
}
200200
Ok(resp)
201201
}
202+
203+
/// REQ-290① 流式精修(NDJSON 逐节,默认路径):SSE 增量累积 → 行级解析
204+
/// (每行=与整包数组元素同构的节对象);每解析一节即回调 on_section
205+
/// (调用方渲染推 Delta 帧——打字机正文)。返回整包 Response(schema v2,
206+
/// 与终稿同构;下游 diff/落库语义零变化)。
207+
///
208+
/// @ai-context: 模型不遵守逐节约定(整包 JSON 输出)→ sections 空 → Err
209+
/// Parse,调用方同 attempt 回退非流式 refine()(行为与旧版
210+
/// 逐字节一致——流式只是呈现增强,诚实降级)。
211+
#[allow(clippy::type_complexity)]
212+
pub fn refine_stream_ndjson(
213+
&self,
214+
request: &AiRefineRequest,
215+
images: &[String],
216+
dims: Option<&ResolvedDims>,
217+
mut on_section: impl FnMut(crate::ai_refine_protocol::AiRefineSection),
218+
) -> Result<AiRefineResponse, AiClientError> {
219+
use crate::ai_refine_protocol::parse_section_ndjson_line;
220+
221+
let mut system = self.prompt.build_system(&request.profile, dims);
222+
system.push_str("\n\n");
223+
system.push_str(crate::ai_refine_protocol::NDJSON_SYSTEM_SUFFIX);
224+
// REQ-290② 预算(与 refine_vision 同口径——流式不豁免上限)
225+
let budget = crate::refine_budget::output_budget(
226+
dims.map(|d| d.preset_id.as_str()).unwrap_or("standard"),
227+
request.content.chars().count(),
228+
);
229+
let user = serde_json::to_string(request)
230+
.map_err(|e| AiClientError::Parse(format!("精修请求序列化失败: {}", e)))?;
231+
let payload = if images.is_empty() {
232+
crate::ai_client::build_chat_payload(
233+
&self.client.config.model, &system, &user, budget.max_tokens,
234+
)
235+
} else {
236+
crate::ai_client::build_vision_payload(
237+
&self.client.config.model, &system, &user, images, budget.max_tokens,
238+
)
239+
};
240+
let mut sections: Vec<crate::ai_refine_protocol::AiRefineSection> = Vec::new();
241+
let mut pending = String::new();
242+
let outcome = crate::ai_chat_stream::stream_sse_content(
243+
&self.client,
244+
payload,
245+
|delta| {
246+
pending.push_str(delta);
247+
while let Some(pos) = pending.find('\n') {
248+
let line: String = pending.drain(..=pos).collect();
249+
if let Some(sec) = parse_section_ndjson_line(&line) {
250+
sections.push(sec.clone());
251+
on_section(sec);
252+
}
253+
}
254+
},
255+
)?;
256+
// 末行无换行(SSE 结束前 flush)
257+
if let Some(sec) = parse_section_ndjson_line(&pending) {
258+
sections.push(sec.clone());
259+
on_section(sec);
260+
}
261+
let _ = outcome.content;
262+
if sections.is_empty() {
263+
return Err(AiClientError::Parse(
264+
"流式响应未解析出章节(模型未按逐节输出——回退非流式)".to_string(),
265+
));
266+
}
267+
Ok(crate::ai_refine_protocol::AiRefineResponse {
268+
schema_version: crate::ai_refine_protocol::SCHEMA_VERSION_V2,
269+
sections,
270+
})
271+
}
202272
}
203273

204274
/// 单测独立文件(保持本文件 ≤300 行,AGENTS.md §3)。

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

Lines changed: 83 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -193,45 +193,70 @@ impl AiRefineResponse {
193193
/// image=原样配图行 `- ![画面](path)`(F3 v2——与规则版
194194
/// 画面要点行同形态,前端渲染器统一处理)。
195195
pub fn to_markdown(&self) -> String {
196-
let mut out = String::new();
197-
for sec in &self.sections {
198-
out.push_str(&format!("## {}\n\n", sec.heading.trim()));
199-
for b in &sec.blocks {
200-
match b.block_type {
201-
AiRefineBlockType::Paragraph => {
202-
out.push_str(b.content.trim());
203-
out.push_str("\n\n");
204-
}
205-
AiRefineBlockType::List => {
206-
for line in b.content.lines().map(|l| l.trim()).filter(|l| !l.is_empty()) {
207-
out.push_str(&format!("- {}\n", line));
208-
}
209-
out.push('\n');
210-
}
211-
AiRefineBlockType::Term => {
212-
out.push_str(&format!("- **{}**\n", b.content.trim()));
213-
}
214-
AiRefineBlockType::Highlight => {
215-
out.push_str(&format!("**{}**\n\n", b.content.trim()));
216-
}
217-
AiRefineBlockType::Quote => {
218-
for line in b.content.lines().map(|l| l.trim()).filter(|l| !l.is_empty()) {
219-
out.push_str(&format!("> {}\n", line));
220-
}
221-
out.push('\n');
196+
render_sections(&self.sections).trim().to_string()
197+
}
198+
}
199+
200+
/// 章节序列 → Markdown(纯函数;to_markdown 与流式 Delta 渲染共用同一出口——
201+
/// REQ-290① 增量帧按前缀渲染,文本与终稿逐字节一致)。
202+
pub fn render_sections(sections: &[AiRefineSection]) -> String {
203+
let mut out = String::new();
204+
for sec in sections {
205+
out.push_str(&format!("## {}\n\n", sec.heading.trim()));
206+
for b in &sec.blocks {
207+
match b.block_type {
208+
AiRefineBlockType::Paragraph => {
209+
out.push_str(b.content.trim());
210+
out.push_str("\n\n");
211+
}
212+
AiRefineBlockType::List => {
213+
for line in b.content.lines().map(|l| l.trim()).filter(|l| !l.is_empty()) {
214+
out.push_str(&format!("- {}\n", line));
222215
}
223-
AiRefineBlockType::Image => {
224-
// F3 v2:配图行(content=session-images/{sid}/{rel}——
225-
// 原样输出;与规则版画面要点行同形态供前端渲染)
226-
out.push_str(&format!("- ![画面]({})\n", b.content.trim()));
216+
out.push('\n');
217+
}
218+
AiRefineBlockType::Term => {
219+
out.push_str(&format!("- **{}**\n", b.content.trim()));
220+
}
221+
AiRefineBlockType::Highlight => {
222+
out.push_str(&format!("**{}**\n\n", b.content.trim()));
223+
}
224+
AiRefineBlockType::Quote => {
225+
for line in b.content.lines().map(|l| l.trim()).filter(|l| !l.is_empty()) {
226+
out.push_str(&format!("> {}\n", line));
227227
}
228+
out.push('\n');
229+
}
230+
AiRefineBlockType::Image => {
231+
// F3 v2:配图行(content=session-images/{sid}/{rel}——
232+
// 原样输出;与规则版画面要点行同形态供前端渲染)
233+
out.push_str(&format!("- ![画面]({})\n", b.content.trim()));
228234
}
229235
}
230236
}
231-
out.trim().to_string()
232237
}
238+
out
233239
}
234240

241+
/// 单行 NDJSON 节解析(REQ-290① 流式):每节=与终稿数组元素同构的 JSON 对象
242+
/// (heading + blocks);非对象行/解析失败 → None(由上层按整包回退解析)。
243+
pub fn parse_section_ndjson_line(line: &str) -> Option<AiRefineSection> {
244+
let t = line.trim();
245+
if t.is_empty() || t.starts_with('[') || t.starts_with(']') || t.starts_with(',') {
246+
return None;
247+
}
248+
let v: serde_json::Value = serde_json::from_str(t).ok()?;
249+
let obj = v.as_object()?;
250+
if !obj.contains_key("heading") || !obj.contains_key("blocks") {
251+
return None;
252+
}
253+
serde_json::from_value::<AiRefineSection>(v).ok()
254+
}
255+
256+
/// 流式模式输出要求(附加在 system 尾部——逐节 JSONL;无法遵守时整体 JSON
257+
/// 数组亦可——解析端自动兼容回退,见 adapter)。
258+
pub const NDJSON_SYSTEM_SUFFIX: &str = "【流式输出要求】请逐节输出整理结果:每一节单独输出一行 JSON(对象含 heading 与 blocks 字段,字段格式与前述格式要求完全一致),节与节之间只换行、不写数组括号/逗号/多余说明;若必须一次性输出,则输出完整 JSON 数组(系统可自动兼容解析)。";
259+
235260
/// 单测独立文件(保持本文件 ≤300 行,AGENTS.md §3)。
236261
#[cfg(test)]
237262
#[path = "ai_refine_protocol_tests.rs"]
@@ -241,3 +266,30 @@ mod tests;
241266
#[cfg(test)]
242267
#[path = "refine_golden_tests.rs"]
243268
mod golden_tests;
269+
270+
#[cfg(test)]
271+
mod ndjson_parse_tests {
272+
use super::parse_section_ndjson_line;
273+
274+
#[test]
275+
fn parses_single_section_object_line() {
276+
let line = r#"{"heading":"第一节","blocks":[]}"#;
277+
let sec = parse_section_ndjson_line(line).expect("应解析出节");
278+
assert_eq!(sec.heading, "第一节");
279+
assert!(sec.blocks.is_empty());
280+
}
281+
282+
#[test]
283+
fn ignores_array_wrappers_and_empty_lines() {
284+
assert!(parse_section_ndjson_line("[").is_none());
285+
assert!(parse_section_ndjson_line("]").is_none());
286+
assert!(parse_section_ndjson_line(",").is_none());
287+
assert!(parse_section_ndjson_line(" ").is_none());
288+
}
289+
290+
#[test]
291+
fn ignores_non_section_objects_and_garbage() {
292+
assert!(parse_section_ndjson_line(r#"{"a":1}"#).is_none());
293+
assert!(parse_section_ndjson_line("not json {").is_none());
294+
}
295+
}

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

Lines changed: 50 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,9 @@ pub enum RefineStreamFrame {
4545
Progress { slice_index: usize, slice_total: usize },
4646
/// 该片精修结果(validate 通过后的渲染 markdown——逐章正文流出)
4747
BlockDone { slice_index: usize, markdown: String },
48+
/// REQ-290①:片内流式增量(NDJSON 逐节解析即推——打字机正文;text=节渲染
49+
/// 的 markdown 片段,同片内按到达序拼接;终稿以 BlockDone 为准)
50+
Delta { slice_index: usize, text: String },
4851
/// 该片失败(回退纯规则语义——诚实降级提示)
4952
SliceFailed { slice_index: usize, reason: String },
5053
/// 任务终态(全部片合并完成)
@@ -510,6 +513,7 @@ pub(crate) fn refine_slices_concurrent(ctx: RefineCtx<'_>) -> (Vec<String>, usiz
510513
let mock_adapter = ctx.mock_adapter;
511514
let mock = ctx.mock;
512515
let task_id = ctx.task_id;
516+
let st = ctx.st; // REQ-290①:流式 Delta 帧推送(worker 线程内 emit)
513517
let profile = ctx.profile; // 轨迹 system 提示词构建用(模板分组同请求)
514518
let vision_images = ctx.vision_images;
515519
let dims = ctx.dims; // 策略解析结果(worker 只读共享)
@@ -524,11 +528,55 @@ pub(crate) fn refine_slices_concurrent(ctx: RefineCtx<'_>) -> (Vec<String>, usiz
524528
let mut outcome: Option<String> = None;
525529
for attempt in 0..=SLICE_RETRY {
526530
// REQ-290(v0.19.6)埋点先行:单片耗时归因——任务级 elapsed_ms
527-
// 已有(db_ai_tasks),此处补片级时间(非流式阶段只有总耗时;
528-
// 流式上线后同点补首 delta 时刻,见批次设计 §2.8 ③)。
531+
// 已有(db_ai_tasks),此处补片级时间(含流式整包耗时口径)。
529532
let started = std::time::Instant::now();
533+
// REQ-290①(v0.19.7):首拍(attempt 0)优先流式(NDJSON 逐节,
534+
// 解析一节推一节 Delta——打字机正文);流式不可用/模型未遵守
535+
// 逐节约定 → 同拍回退非流式(与旧版逐字节一致);重试拍走非
536+
// 流式(幂等语义下不重复推流)。
537+
let stream_enabled = !mock
538+
&& attempt == 0
539+
&& std::env::var("REFINE_STREAM_NDJSON")
540+
.map(|v| v != "0")
541+
.unwrap_or(true);
530542
let resp = if mock {
531543
Ok(mock_adapter.refine(req))
544+
} else if stream_enabled {
545+
match adapter
546+
.refine_stream_ndjson(
547+
req,
548+
ctx.vision_images,
549+
Some(dims),
550+
|sec| {
551+
emit_refine_stream(
552+
st,
553+
task_id,
554+
RefineStreamFrame::Delta {
555+
slice_index: idx + 1,
556+
text: crate::ai_refine_protocol::render_sections(
557+
std::slice::from_ref(&sec),
558+
),
559+
},
560+
);
561+
},
562+
)
563+
.map_err(AiTaskFailure::from)
564+
{
565+
Ok(r) => Ok(r),
566+
Err(e) => {
567+
eprintln!(
568+
"[refine-task] task={} 片 {} 流式不可用,同拍回退非流式: {}",
569+
task_id, idx + 1, e.message()
570+
);
571+
if ctx.vision_images.is_empty() {
572+
adapter.refine(req, Some(dims)).map_err(AiTaskFailure::from)
573+
} else {
574+
adapter
575+
.refine_vision(req, ctx.vision_images, Some(dims))
576+
.map_err(AiTaskFailure::from)
577+
}
578+
}
579+
}
532580
} else if ctx.vision_images.is_empty() {
533581
adapter.refine(req, Some(dims)).map_err(AiTaskFailure::from)
534582
} else {

‎app/src/components/TaskThreadCard.test.tsx‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@ import type { AiTaskRecord } from "../types";
1414
// v0.17.0:流式 hook 依赖 tauri event listen——测试环境 mock 为空流(不报错)
1515
vi.mock("../hooks/useRefineStream", () => ({
1616
useRefineStream: () => [],
17-
orderedBlockFrames: () => [],
17+
sliceStreamContent: () => [],
1818
}));
1919

2020
function task(partial: Partial<AiTaskRecord> & { taskId: number }): AiTaskRecord {

‎app/src/components/TaskThreadCard.tsx‎

Lines changed: 12 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,8 @@
1010
* [回到会话] [查看笔记](用户裁决:完成后快速闭环)。
1111
*/
1212
import type { AiTaskRecord } from "../types";
13-
import { orderedBlockFrames, useRefineStream } from "../hooks/useRefineStream";
13+
// REQ-290①(v0.19.7):delta 逐节增量并入片正文(打字机);blockDone 收尾
14+
import { sliceStreamContent, useRefineStream } from "../hooks/useRefineStream";
1415

1516
interface Props {
1617
tasks: AiTaskRecord[];
@@ -51,26 +52,29 @@ function recentDone(tasks: AiTaskRecord[]): AiTaskRecord | null {
5152
/** 精修流式正文(进行中任务卡内——逐章流出;无帧则仅进度行) */
5253
function RefineStreamBody({ taskId, total }: { taskId: number; total: number | null }) {
5354
const frames = useRefineStream(taskId);
54-
const blocks = orderedBlockFrames(frames);
55+
const slices = sliceStreamContent(frames);
5556
const doneFrames = frames.filter((f) => f.kind === "done").length;
5657
const failedCount = frames.filter((f) => f.kind === "sliceFailed").length;
57-
if (blocks.length === 0 && frames.length === 0) return null;
58+
const streamingLive = slices.some((s) => !s.complete && s.text.length > 0);
59+
if (slices.length === 0 && frames.length === 0) return null;
5860
return (
5961
<div style={{ width: "100%", marginTop: 4, borderTop: "1px dashed #e5e7eb", paddingTop: 4 }}>
6062
<div style={{ fontSize: 10.5, color: "#9ca3af", marginBottom: 2 }}>
61-
已整理 {blocks.length}/{total ?? "?"} 片{failedCount > 0 && ` · 失败 ${failedCount} 片(保留规则版)`}
63+
已整理 {slices.filter((s) => s.complete).length}/{total ?? "?"} 片{failedCount > 0 && ` · 失败 ${failedCount} 片(保留规则版)`}
64+
{streamingLive && " · 逐节流出中"}
6265
{doneFrames > 0 && " · 完成"}
6366
</div>
64-
{blocks.map((b) => (
67+
{slices.map((s) => (
6568
<pre
66-
key={b.sliceIndex}
69+
key={s.sliceIndex}
6770
style={{
6871
whiteSpace: "pre-wrap", wordBreak: "break-word", fontFamily: "inherit",
6972
fontSize: 11.5, color: "#374151", margin: 0, padding: "2px 0",
70-
borderTop: b.sliceIndex > 1 ? "1px solid #f3f4f6" : "none",
73+
borderTop: s.sliceIndex > 1 ? "1px solid #f3f4f6" : "none",
7174
}}
7275
>
73-
{b.markdown}
76+
{s.text}
77+
{!s.complete && s.text.length > 0 ? " ▍" : ""}
7478
</pre>
7579
))}
7680
</div>

0 commit comments

Comments
 (0)