AI 回复为什么要在 Rust 里收流:Nexora 的 FRB 流式对话实战
项目仓库:https://atomgit.com/nutpi/Nexora
FRB 项目:https://atomgit.com/oh-flutter/flutter_rust_bridge
Rust 社区:https://xuanwu.openatom.cn/
聊天应用最容易写出的版本,是 Dart 发一个 HTTP 请求,等完整文本回来后塞进消息气泡。真正使用时,这个方案的问题很快就出现了:首字等待太久、用户无法取消、页面切走后请求还在跑、流到一半退出会丢掉已经生成的内容。
Nexora 是一个面向 HarmonyOS PC 的 AI 工作区。它把供应商协议、网络流、取消控制和 SQLite 持久化放在 Rust,Flutter 只消费统一的事件流。FRB 在这里不是普通的函数调用通道,而是把 Rust 的异步过程变成 Dart Stream。
本文按“认识 FRB 与 Rust 工具链、理解三层职责、阅读核心代码、完成真机验证、定位问题和参与项目”的顺序展开。已有章节保留原有的流式实现和验收细节,新增内容只用于补充项目入口与复现路径。
技术选型与复现入口
1. Flutter、Rust 与 FRB 的职责
| 组成 | 在 Nexora 中的职责 | 选择它的价值 |
|---|---|---|
| Flutter | 工作区、消息列表、输入框和流式渲染 | 负责交互和展示,不持有网络任务的最终状态 |
| Rust | 供应商协议、HTTP 流、取消、SQLite 和消息状态 | 让长期运行的任务、部分文本和终态落在同一个边界 |
| FRB | 生成 Dart/Rust 绑定,传递事件、结果和错误 | 以 Stream<ChatStreamEvent> 暴露稳定业务语义 |
2. 工具链检查
复现前先确认 Flutter-OH、Rust 和 FRB 生成器版本与仓库构建说明一致:
flutter --version
flutter doctor -v
rustc --version
cargo --version
rustup target list --installed
flutter pub get
flutter_rust_bridge_codegen generate
生成绑定后再构建 HarmonyOS 应用。若问题只在最小 FRB 示例中出现,优先到 FRB 仓库提交 Issue;若只在消息状态、供应商解析或 SQLite 中出现,则先在 Nexora 仓库定位。
3. 旋武社区与 FRB
开放原子旋武开源社区提供 Rust 学习、工具和开源协作入口;flutter_rust_bridge负责生成 Dart/Rust 跨语言绑定。Nexora 的业务问题进入项目仓库,能在最小工程复现的绑定或类型映射问题再进入 FRB 仓库。
最终界面

图中是 Nexora 在 HarmonyOS PC 真机中的实际运行状态。当前会话使用 deepseek-chat | deepseek,用户要求精确回复 DEEPSEEK_OK,服务端返回了对应文本;随后发送 hi,也收到了完整回复。截图只展示会话结果,没有打开设置页,因此不会暴露 API Key。
此前的内存仓库截图已经撤下。内存仓库适合覆盖加载、空态和错误态,却不能证明网络请求、Rust 流解析或 SQLite 持久化真的工作。
先看目录,再看流
Nexora/
|-- lib/main.dart RustLib 初始化
|-- lib/src/chat/chat_repository.dart Dart 侧聊天契约
|-- lib/src/workspace/workspace_page.dart 三栏工作区
|-- lib/src/providers/ 服务商设置和仓库
|-- rust/src/api/chat.rs 会话与消息 CRUD
|-- rust/src/api/ai_stream.rs 流请求、取消和事件
|-- rust/src/api/provider_profile.rs 模型与服务商配置
|-- rust/src/storage/ SQLite、迁移和查询
`-- lib/src/rust/ FRB 生成绑定
排查“为什么没回复”时,先从消息记录找 assistantMessage.id,再查它是否进入活动流表,最后看供应商解析器。不要一开始就盯着消息气泡;页面只是最终消费者,真正的任务身份在 Rust 和数据库里。
一条消息实际经过哪些层
这里特意把“创建消息”和“开始收流”拆开。messageId 是一次生成任务的稳定身份,取消、断点保存和错误恢复都围绕它进行,而不是围绕某个 Widget 的生命周期。
FRB 配置和初始化
rust_input: crate::api
rust_root: rust/
dart_output: lib/src/rust
应用启动时先初始化原生库,再注入生产仓库:
Future<void> main() async {
WidgetsFlutterBinding.ensureInitialized();
await RustLib.init();
runApp(
NexoraApp(
chatRepository: const NativeChatRepository(),
// 其他仓库省略
),
);
}
页面不直接 import 一堆 FRB 函数,而是使用 ChatRepository:
abstract interface class ChatRepository {
Future<native.ChatTurnSnapshot> createTurn(
native.CreateChatTurnRequest request,
);
Stream<ai.ChatStreamEvent> streamResponse(String messageId);
Future<bool> cancelResponse(String messageId);
}
class NativeChatRepository implements ChatRepository {
const NativeChatRepository();
Stream<ai.ChatStreamEvent> streamResponse(String messageId) =>
ai.streamChatResponse(messageId: messageId);
Future<bool> cancelResponse(String messageId) =>
ai.cancelChatStream(messageId: messageId);
}
这样做不只是为了测试。以后网络层迁移、增加本地模型或回放历史流时,页面依然只认识同一个仓库契约。
本项目使用的 Rust 特性
Nexora 中的 Rust 特性都直接服务于“流任务可取消、可恢复、可持久化”这条链路:
| Rust 特性 | 在 Nexora 中的实际用法 |
|---|---|
struct 与 enum | 用 ChatStreamEvent 承载事件,用 ChatStreamEventKind 枚举 Delta、Completed、Failed、Paused 等有限状态;派生 Debug、Clone、PartialEq 便于日志和测试。 |
Option<T> | 错误码和错误消息允许缺省,避免用空字符串区分“没有错误”和“错误内容为空”。 |
Result<T, E> 与 ? | 网络、数据库和状态注册失败沿调用链返回 NativeError,不会被静默吞掉。 |
async/await | stream_chat_response 以异步函数读取 HTTP 流,在等待网络和数据库时不阻塞 Flutter 线程。 |
模式匹配 match | 按 ReceiveResult 的完成、暂停、失败分支执行不同的持久化和事件发布动作。 |
| 所有权与借用 | &context、&provider、&mut text 只在收流期间借用数据,消息文本的最终所有权在终态保存时明确转移。 |
RAII 与 Drop | ActiveStreamGuard 离开作用域时自动清理活动任务表,即使中途返回错误也不会遗留任务。 |
Tokio watch 通道 | 用共享的取消信号通知网络读取任务;Flutter 只传 message_id,不暴露 Rust task 句柄。 |
FRB 的 StreamSink | 通过 StreamSink<ChatStreamEvent> 把 Rust 事件流安全地映射成 Dart Stream。 |
Rust 怎样把异步响应送回 Flutter
跨语言事件定义保持得很克制:
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ChatStreamEventKind {
Delta,
Completed,
Failed,
Paused,
}
#[derive(Debug, Clone)]
pub struct ChatStreamEvent {
pub message_id: String,
pub kind: ChatStreamEventKind,
pub delta: String,
pub text: String,
pub error_code: Option<String>,
pub error_message: Option<String>,
}
delta 适合即时追加,text 是截至当前事件的完整文本。只传 delta 会让重连和状态校正变麻烦;每次只传完整文本又会增加复制。两者同时保留,UI 可以快速渲染,终态也能校验。
核心入口直接接收 StreamSink:
pub async fn stream_chat_response(
message_id: String,
sink: StreamSink<ChatStreamEvent>,
) -> Result<(), NativeError> {
let context = load_chat_generation_context(&message_id)?;
let provider = resolve_chat_provider(&context.model_record_id)?;
let mut cancellation = register_stream(&message_id)?;
let _guard = ActiveStreamGuard::new(message_id.clone());
let mut text = String::new();
match receive_response(
&message_id,
&context,
&provider,
&sink,
&mut cancellation,
&mut text,
PersistTarget::Chat,
).await {
Ok(ReceiveResult::Completed) => {
persist_assistant_stream_state(
&message_id,
&text,
ChatMessageStatus::Success,
)?;
let _ = sink.add(ChatStreamEvent {
message_id,
kind: ChatStreamEventKind::Completed,
delta: String::new(),
text,
error_code: None,
error_message: None,
});
Ok(())
}
// Paused 和 Failed 分支省略
}
}
ActiveStreamGuard 利用 Drop 清理活动任务表,避免异常返回后留下“该消息仍在生成”的假状态。取消信号用 Tokio watch 通道传递;Dart 调 cancelChatStream 时只需要消息 ID,不需要持有 Rust task 句柄。
多家模型协议怎样收敛
Nexora 当前不是只支持一种 SSE。Rust 内部将供应商分为 OpenAI Chat Completions、OpenAI Responses、Anthropic Messages、Google 和 Ollama,再把各自的帧解析成统一的 ParsedFrame { delta, done, provider_error }。
供应商差异留在 Rust 的理由很直接:
- SSE 分片可能把一段 JSON 拆成多次网络读取;
- 不同供应商的结束标记和错误载荷不同;
- API Key、Endpoint 和模型覆盖配置都来自本地存储;
- 同一套取消和部分保存策略应该对所有协议生效。
代码还设置了 MAX_RESPONSE_CHARS = 200_000,并按 512 字节或 500ms 的间隔保存部分文本。限制响应上限不是保守过头,而是避免异常服务端把桌面进程拖进无限增长。
从发送到取消的实操
Flutter 侧可以按下面的顺序处理一次对话:
final turn = await repository.createTurn(request);
final subscription = repository
.streamResponse(turn.assistantMessage.id)
.listen((event) {
switch (event.kind) {
case ChatStreamEventKind.delta:
state = state.copyWith(text: event.text, generating: true);
case ChatStreamEventKind.completed:
case ChatStreamEventKind.paused:
case ChatStreamEventKind.failed:
state = state.copyWith(text: event.text, generating: false);
}
});
// 用户点击停止
await repository.cancelResponse(turn.assistantMessage.id);
本地验证命令:
cd Nexora
flutter pub get
flutter_rust_bridge_codegen generate
cargo test --manifest-path rust/Cargo.toml
flutter test
flutter run -d macos
测试时至少覆盖完成、主动取消、供应商返回错误、空响应、页面离开这五种路径。尤其要确认取消后的部分文本状态是 Paused,而不是被当成成功或直接丢弃。
配置服务商时别把密钥带进截图
真实联调需要 Endpoint、模型名和凭据,但文章素材不应该出现这些值。我的做法是先在设置页完成配置并保存,回到会话页后再截图。页面只显示服务商类型和模型名,凭据既不进入日志,也不出现在图片中。
如果需要记录请求排查信息,只保留以下字段就够了:provider kind、最终 URL 的域名、模型、HTTP 状态码、错误码、首字耗时和总耗时。Authorization 头、完整请求体中的私有系统提示、代理地址参数都不应写进公开博客。
真机验证步骤
HDC=/Applications/DevEco-Studio.app/Contents/sdk/default/openharmony/toolchains/hdc
BUNDLE_NAME=$(awk -F'"' '/"bundleName"/ {print $4; exit}' ohos/AppScope/app.json5)
$HDC shell aa start -a EntryAbility -b "$BUNDLE_NAME"
$HDC shell uitest dumpLayout -b "$BUNDLE_NAME" \
-p /data/local/tmp/nexora-layout.json
$HDC shell snapshot_display -f /data/local/tmp/nexora-chat.jpeg
这次界面树同时读到了会话标题、deepseek-chat | deepseek、用户消息 Reply with exactly: DEEPSEEK_OK 和助手消息 DEEPSEEK_OK。因此图片对应的是应用真实状态,不是把一张对话图片贴进页面。
流式状态不能只有 loading 和 done
一次生成至少会落在以下状态之一:
| 终态 | 数据库内容 | 页面处理 |
|---|---|---|
Success | 保存最终完整文本 | 结束光标动画,允许复制 |
Paused | 保存取消前的部分文本 | 显示已停止,可继续或重试 |
Failed | 保存部分文本和错误信息 | 不清空已经收到的内容 |
| 空响应 | 保存明确错误,而不是空成功 | 提示供应商未返回正文 |
StreamSubscription.cancel() 只是不再监听 Dart Stream,不等于取消 Rust 中的网络任务。真正的停止动作必须调用 cancelChatStream(messageId);Rust 收到信号后结束网络读取、保存现有文本、发布 Paused,最后由 ActiveStreamGuard 清理注册表。这四步少一步,下次生成都可能收到“任务已存在”。
用故障注入检查解析器
流式 HTTP 并不保证一次网络 chunk 就是一条完整 SSE。建议准备一个本地测试服务,刻意把 JSON 从中间切开,再插入空行、注释行和多字节 UTF-8 边界。解析器应先按 SSE 事件聚合,再解释 JSON,不能直接对每次 bytes_stream 回调做 serde_json::from_slice。
还应覆盖下面几种返回:
HTTP 401:凭据错误,不能无限重试
HTTP 429:保留供应商错误和 retry-after
HTTP 200 + provider error frame:按失败终止
正常 delta 后连接断开:保存 partial text 为 failed/paused
收到 done 后还有尾部空行:只发布一次 completed
输出超过 200000 字符:主动停止并返回上限错误
这些用例比“模型能回复你好”更能验证 Rust 收流层的价值。
一份可交付的验收表
| 检查项 | 本次状态 | 判断依据 |
|---|---|---|
| HarmonyOS 真机启动 | PASS | 应用包实际进入三栏工作区 |
| 线上模型回复 | PASS | DeepSeek 按要求返回 DEEPSEEK_OK |
| 第二轮连续对话 | PASS | hi 获得完整助手回复 |
| 会话历史恢复 | PASS | 启动后列表和消息仍存在 |
| 主动取消 | 代码已实现,交付时专项复测 | 需要长回复场景记录 Paused |
| 多供应商兼容 | 解析器和测试覆盖 | 不由单张 DeepSeek 截图证明 |
| 密钥未泄露 | PASS | 截图和正文均未包含凭据 |
流式回复真正改善的不是动画
用户看到文字一段段出现,很容易把流式对话理解成一种视觉效果。对工程实现来说,更重要的是它改变了失败的含义。非流式请求在最后一刻断开,通常只能得到一次失败;流式请求即使中途断开,也可能已经收到一段有价值的内容。系统需要保存这段内容、标记为什么停止,并允许用户决定继续、重试还是保留,而不是把消息气泡清空。
首字时间和总完成时间也应该分开。一个模型可能很快返回第一个 token,但完整回答需要很久;另一个模型首字稍慢,后续吞吐却更稳定。只记录总耗时无法解释用户为什么觉得“卡”。Rust 收流层天然知道请求发起、首个 delta、最后事件和取消时间,可以生成更准确的性能记录,Flutter 不需要自己猜测网络阶段。
流式展示还会放大渲染问题。如果每个极小 delta 都触发整个消息列表重建,长回答时 CPU 和内存占用会明显增加。UI 可以在内存中追加文本,并按帧或很短时间窗口合并刷新;Rust 侧仍按自己的字节或时间阈值持久化。渲染节流和磁盘保存不是同一个频率,不能为了界面流畅而减少可靠保存,也不能每收到一个字符就写一次数据库。
生成任务为什么必须属于消息,而不是页面
页面会被重建、切换甚至销毁,消息 ID 却可以稳定存在。把任务绑在消息 ID 上,用户离开当前话题时,Rust 仍能完成生成并保存结果;回到话题后,页面从数据库读取最新状态即可。若业务规则要求离开页面就停止,也应显式调用取消,而不是期待 Widget dispose 顺便终止底层网络连接。
这个所有权模型对重复发送也有帮助。同一个 assistant message 只能有一个活动任务,Rust 注册表可以拒绝第二次启动。否则用户快速按两次发送,两个响应同时向同一个气泡追加,数据库里会出现无法解释的混合文本。前端禁用按钮只能减少误触,真正的唯一性仍要在任务管理层保证。
应用异常退出后,数据库中可能留有 pending 消息。下一次启动不能永久显示正在生成,因为原任务已经不存在。比较合理的恢复方式是把没有活动任务的 pending 状态转换为 paused 或 interrupted,保留已经写入的部分文本。用户看到中断原因后可以重新生成,而不是面对一个永远转圈的会话。
多供应商统一最难的是错误语义
正常文本 delta 通常都能映射成字符串,真正麻烦的是错误和终止。有的服务用 HTTP 状态码返回错误,有的先返回 200,再在 SSE JSON 里放错误对象;有的发送明确 done 标记,有的直接关闭连接;还有的把推理内容、工具调用和最终文本分在不同字段。统一接口如果只抽象正常文本,页面仍会充满供应商判断。
Nexora 把这些差异收进 Rust 后,Flutter 只处理 Delta、Completed、Paused 和 Failed。这个统一并不意味着丢掉细节。Failed 事件仍要保留稳定错误码、适合用户阅读的信息和用于日志定位的供应商上下文。对 401、429、上下文超限和模型不存在这几类常见错误,页面给出的下一步应不同:检查凭据、稍后重试、缩短上下文或重新选择模型。
结束事件还必须保证只有一次。某些协议同时包含内容结束字段和流关闭,若两个路径都发布 Completed,页面可能执行两次保存、滚动或通知。Rust 内部应在状态转换时锁定终态,后续尾部空行和连接关闭只负责清理。取消与自然完成在临界点竞争时,也要保证数据库最终只有一个状态,而不是先 Success 后又被 Paused 覆盖。
真实 DeepSeek 会话能证明到哪一步
这次真机截图中,第一轮要求模型只回复 DEEPSEEK_OK,最终内容与要求一致。这样做不是为了展示模型有多聪明,而是便于核对请求和响应有没有串话、前后有没有夹带隐藏提示文本。第二轮发送简短的 hi,确认同一话题可以继续创建新的 user/assistant 消息。应用重启后仍能看到两轮内容,说明历史读取也正常工作。
不过这张图看不到流式过程中每一个 delta,也没有展示主动取消,所以文章没有把这两项写成由截图直接证明。流解析能力由 Rust 测试覆盖,取消则还应在交付环境发起一条足够长的回答,在中途停止,核对页面状态、部分文本和数据库记录。把证据与结论一一对应,比在一张最终截图上解释所有功能更可靠。
截图还刻意避开了设置页。真实 API 调用需要密钥,但验收图片的目的不是证明密钥存在。设置页一旦出现完整 Endpoint 查询参数、代理认证或可见 Token,就会把一次功能证明变成安全事故。最稳妥的做法是配置完成后回到会话页,只保留模型和供应商名称,并在发布前再做一次图像人工检查。
数据库设计不能只考虑最终消息
一次对话通常先创建用户消息和空的助手消息,再不断更新助手内容。若两条记录不在同一事务里,模型配置无效时可能只留下用户消息;若更新没有检查状态,过期任务可能覆盖新一轮重试。消息表至少需要稳定 ID、父子关系、状态和更新时间,流任务则始终根据同一个 message ID 更新。
部分保存的间隔要在可靠性和写放大之间取平衡。每 512 字节或约 500ms 保存一次,意味着崩溃时最多损失一个很短窗口,又不会每个 token 都提交 SQLite。这个数值不是固定真理,大模型输出速度、设备磁盘和消息长度不同,都可能需要调整。重要的是同时设置字节和时间条件:输出很慢时靠时间落盘,输出很快时靠字节阈值限制内存积累。
搜索索引也应只收录适合展示的最终文本。供应商可能返回隐藏推理或内部字段,如果解析时把它们混进消息正文,不仅界面会泄露不该显示的信息,全文检索也会长期保存。Nexora 的测试专门确认文本提取忽略隐藏推理,这类测试看起来不如聊天截图直观,实际上与数据安全直接相关。
做完这个模块后的一个判断
最开始很容易把 FRB 当作“Dart 不方便写网络,所以交给 Rust”。真正做完后会发现,语言不是重点,任务归属才是重点。只要网络任务、取消令牌、部分文本和数据库状态由同一层维护,页面切换和供应商差异就不再破坏一致性。Rust 恰好适合承载这套长期运行的异步状态,而 FRB 提供了清晰的事件出口。
后续如果增加工具调用、图片输入或本地模型,也不应该绕开这个任务模型。它们可以增加新的事件内容,却仍要围绕消息 ID、终态和持久化展开。保持这条主线,功能越多时系统反而越容易解释;一旦每种模型各自管理页面状态,短期开发快,后期排查会迅速失控。
当前状态和限制
网络恢复不能等同于自动重发
对话过程中遇到短暂断网,最危险的处理并不是报错,而是不加判断地自动重发。服务端可能已经收到第一次请求并开始计费,只是客户端没有收到后续流;此时重新提交会产生两条答案,甚至让带副作用的工具调用执行两次。因此 Failed 事件除了错误信息,还应标明请求是否已经发出、是否收到过内容,以及这类请求能否安全重试。
页面上的“重试”也应该创建一条有来源关系的新任务,保留失败记录,而不是覆盖旧消息。这样用户能看见第一次在哪个位置中断,数据库也能解释为什么同一个问题出现两条助手消息。若以后加入服务端幂等键,可以把 message ID 作为请求身份的一部分;在供应商不支持幂等的情况下,则明确采用人工确认重试,更符合真实风险。
恢复网络后,Rust 可以继续处理尚未结束且连接仍有效的流,但不能仅凭系统网络状态从离线变为在线就猜测原请求可恢复。SSE 连接是否支持续传、供应商是否返回事件 ID,都必须以协议事实为准。没有续传能力时,把已有部分文本保存为 Paused 或 Failed,允许用户选择重试,是比拼接两次答案更可解释的行为。
当前项目已经实现会话和消息持久化、服务商配置、统一流式事件、取消、部分文本保存以及多协议解析。本文截图证明本次真机环境中的 DeepSeek 请求与历史恢复已经跑通,但不代表读者环境天然具备同一网络和账号条件;复现时仍需配置自己的合法服务地址、模型和凭据。
这套 FRB 设计真正解决的是“生成任务归谁管”。只要生成任务属于 Rust 和消息 ID,而不是某个临时 Widget,窗口刷新、页面切换、取消与恢复就不再互相打架。
FAQ:从问题定位到 Issue 与 PR
1. 截图能证明 Rust 流式链路吗?
不能。截图只能证明页面最终显示了结果;流解析、主动取消、部分文本保存和终态写入仍需通过 Rust 单测、日志和真机操作分别核对。
2. FRB 生成失败先检查什么?
先确认 Dart 依赖、Rust 依赖和 flutter_rust_bridge_codegen 版本一致,再检查 flutter_rust_bridge.yaml 的 rust_input、rust_root 和 dart_output。最小复现仍失败时,再将版本、平台、完整错误和最小代码提交到 FRB 仓库。
3. 如何提交 Nexora 的修复?
先在本地补充对应测试,记录请求是否发出、是否收到 delta、最终状态和数据库结果,再提交 PR。PR 描述应包含问题触发条件、根因、修改内容、测试命令和 HarmonyOS 真机结果。
后续扩展与项目入口
Nexora 后续可以继续增加本地模型、工具调用和更细的重连策略,但每种能力都应沿用消息 ID、终态和持久化契约。项目源码、FRB 入口和 Rust 社区链接已在文章开头列出,反馈时请同时注明仓库版本和目标平台。
更多推荐



所有评论(0)