Skip to main content

streaming_reader/
streaming_reader.rs

1//! 流式读取:从 `std::io::Read` 增量解码,边到边处理(示例 5/6)。
2//!
3//! 运行:`cargo run -p nextjson --example streaming_reader`
4//!
5//! 网络/管道场景里整段载荷往往不会一次性到齐。`StreamDecoder` 只按需向
6//! `Read` 拉取字节,因此解码可以在数据尚未全部到达时就开始:
7//! - `nextjson::from_reader`:顶层便捷入口(要求拥有型目标);
8//! - `nextjson::StreamDecoder`:手动控制解码进度。
9//!
10//! 示例用一个"每次最多吐 3 字节"的 `ChunkedReader` 模拟慢速 socket,
11//! 证明解码不依赖整包到达。
12
13use std::io::{Cursor, Read};
14
15use nextjson::{from_reader, NsonDeserialize, NsonSerialize, StreamDecoder};
16
17/// 拥有型消息(流式解码要求 `for<'de> NsonDeserialize<'de>`,不能借用输入)。
18#[derive(Debug, PartialEq, NsonSerialize, NsonDeserialize)]
19struct Message {
20    id: u64,
21    body: String,
22}
23
24/// 慢速读取器:每次 `read` 最多返回 `chunk` 字节,模拟分片到达的网络流。
25struct ChunkedReader<R> {
26    inner: R,
27    chunk: usize,
28}
29
30impl<R: Read> Read for ChunkedReader<R> {
31    fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
32        let want = buf.len().min(self.chunk);
33        let mut slice = buf;
34        slice = &mut slice[..want];
35        self.inner.read(slice)
36    }
37}
38
39fn main() -> nextjson::Result<()> {
40    let messages = vec![
41        Message {
42            id: 1,
43            body: "first message".into(),
44        },
45        Message {
46            id: 2,
47            body: "second message".into(),
48        },
49        Message {
50            id: 3,
51            body: "third message".into(),
52        },
53    ];
54
55    // 1. 顶层入口:把一条 JSON 载荷放进 Cursor,再包一层慢速读取器。
56    let payload = nextjson::to_vec(&messages)?;
57    println!("序列化载荷: {} 字节", payload.len());
58
59    let slow = ChunkedReader {
60        inner: Cursor::new(payload.clone()),
61        chunk: 3,
62    };
63    let decoded: Vec<Message> = from_reader(slow)?;
64    assert_eq!(decoded, messages);
65    println!(
66        "from_reader + 3 字节/次: 解码 {} 条消息,与源数据一致",
67        decoded.len()
68    );
69
70    // 2. 手动控制:StreamDecoder 边到边读取单条消息。
71    let single = Message {
72        id: 42,
73        body: "streamed".into(),
74    };
75    let one_payload = nextjson::to_vec(&single)?;
76    let slow_one = ChunkedReader {
77        inner: Cursor::new(one_payload),
78        chunk: 2,
79    };
80
81    let mut decoder = StreamDecoder::new(slow_one);
82    let value: Message = Message::nextdecode(&mut decoder)?;
83    decoder.end()?;
84    assert_eq!(value, single);
85    println!("StreamDecoder + 2 字节/次: 解码 msg id = {}", value.id);
86
87    // 3. 流式容器:数组到达时逐条解码(消息没有一次性到齐也没关系)。
88    let mut full = Cursor::new(nextjson::to_vec(&messages)?);
89    let mut d = StreamDecoder::new(&mut full);
90    let streamed: Vec<Message> = Vec::<Message>::nextdecode(&mut d)?;
91    d.end()?;
92    assert_eq!(streamed, messages);
93    println!(
94        "流式解码整个数组: {} 条消息,首条 id = {}",
95        streamed.len(),
96        streamed[0].id
97    );
98
99    Ok(())
100}