1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
//! **tcp** — a real client/server over a reliable byte stream, reading and writing
//! framed messages on one connection, *without* `try_clone()`.
//!
//! A `#[bin]` enum is a self-delimiting RPC (magic-dispatched), so a reader pulls exactly one
//! message off the stream at a time. The read side wraps the socket in a [`BufSource`] (the
//! `std::io::Read` → bnb `Source` adapter that buffers as the decoder needs bytes, so
//! `decode` blocks until a whole message has arrived). The write side just calls
//! `to_bytes()` + `write_all`.
//!
//! **The duplex trick:** `std`'s `&TcpStream` implements *both* `Read` and `Write`, so the read
//! half (a `BufSource<&TcpStream>`) and the write half (`&TcpStream`) are two shared borrows of
//! the *same* socket — no `try_clone()` (which dups the fd), no ownership split. For halves you
//! need to *move* across threads, the equivalent is `Arc<TcpStream>` (what tokio's `into_split`
//! does under the hood).
//!
//! Run with: `cargo run -p bitsandbytes --example tcp`
use bnb::{BufSource, Source, bin};
use std::io::Write;
use std::net::{TcpListener, TcpStream};
use std::thread;
use tracing::info;
/// A tiny self-delimiting RPC: a `b"RPC"` sync prefix, then a one-byte variant `magic`.
#[bin(big, magic = b"RPC")]
#[derive(Debug, Clone, PartialEq, Eq)]
enum Message {
#[bin(magic = 0x01u8)]
Ping { seq: u32 },
#[bin(magic = 0x02u8)]
Pong { seq: u32 },
#[bin(magic = 0x03u8)]
Echo {
#[brw(count_prefix = u8)] // derived, never stored, checked at encode
#[try_str]
text: Vec<u8>,
},
#[bin(magic = 0x04u8)]
Bye,
}
fn echo(text: &str) -> Message {
Message::Echo {
text: text.as_bytes().to_vec(),
}
}
/// The server's reply to one request (`None` = close the connection).
fn reply_to(req: &Message) -> Option<Message> {
match req {
Message::Ping { seq } => Some(Message::Pong { seq: *seq }),
Message::Echo { text } => Some(Message::Echo { text: text.clone() }), // echo it back
Message::Bye => None, // hang up
Message::Pong { .. } => None, // unexpected here
}
}
/// Serve one connection: read framed requests off the stream and write replies, on the *same*
/// socket — read half is a `BufSource<&TcpStream>`, write half is `&TcpStream` (no `try_clone`).
fn serve(stream: TcpStream) -> std::io::Result<()> {
let mut reader = BufSource::new(&stream); // &TcpStream: Read
let mut writer = &stream; // &TcpStream: Write — the same underlying socket
loop {
// Distinguish a clean close from corruption. When the peer hangs up *between*
// messages, the next decode fails with `Incomplete` at exactly the bit position
// where the message would have started — no byte of it ever arrived. An
// `Incomplete` past that position means the connection died mid-message, and any
// other kind (`BadMagic`, `Io`, …) is a framing/transport fault — those are real
// errors and must be surfaced, not silently treated as end-of-stream.
let start = reader.bit_pos();
let req = match Message::decode(&mut reader) {
Ok(req) => req,
Err(e) if e.is_incomplete() && e.at == start => break, // clean close
Err(e) => panic!("server: framing/transport error: {e}"),
};
info!(?req, "server ← request");
match reply_to(&req) {
Some(reply) => {
writer.write_all(&reply.to_bytes().expect("encode reply"))?;
info!(?reply, "server → reply");
}
None => {
info!("server: closing connection");
break;
}
}
}
Ok(())
}
fn hex(bytes: &[u8]) -> String {
bytes
.iter()
.map(|b| format!("{b:02x}"))
.collect::<Vec<_>>()
.join(" ")
}
fn main() -> Result<(), Box<dyn std::error::Error>> {
tracing_subscriber::fmt()
.with_max_level(tracing::Level::INFO)
.with_target(false)
.without_time()
.init();
// A server on an ephemeral loopback port, handling exactly one connection in a thread.
let listener = TcpListener::bind("127.0.0.1:0")?;
let addr = listener.local_addr()?;
let server = thread::spawn(move || {
let (stream, peer) = listener.accept().expect("accept");
info!(%peer, "server: connection accepted");
serve(stream).expect("serve");
});
// The client: one connection, write requests and read the streamed replies — same duplex
// trick (a `BufSource<&TcpStream>` to read, `&TcpStream` to write; no `try_clone`).
let stream = TcpStream::connect(addr)?;
let mut reader = BufSource::new(&stream);
let mut writer = &stream;
for request in [Message::Ping { seq: 1 }, echo("hello"), echo("world")] {
let bytes = request.to_bytes()?;
writer.write_all(&bytes)?;
info!(?request, bytes = %hex(&bytes), "client → request");
let reply = Message::decode(&mut reader)?; // one framed reply off the stream
info!(?reply, "client ← reply");
match (&request, &reply) {
(Message::Ping { seq }, Message::Pong { seq: back }) => assert_eq!(seq, back),
(Message::Echo { text }, Message::Echo { text: back }) => assert_eq!(text, back),
_ => panic!("unexpected reply {reply:?}"),
}
}
// Say goodbye, which tells the server to close.
writer.write_all(&Message::Bye.to_bytes()?)?;
info!("client → Bye");
server.join().expect("server thread panicked");
info!("all checks passed");
Ok(())
}