oxedyne/fe2o3/fe2o3_shield/tests/wire.rs
14.9 KiB, 1 run
created by r1870400018:20748, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | //! Two peers on one machine, and the bytes that get from one to the other and |
| 2 | //! back. |
| 3 | //! |
| 4 | //! Everything here is the real wire: real chunking, a real proof of work on |
| 5 | //! every packet, a real Ed25519 signature over every packet, the real address |
| 6 | //! guard, and the real message assembler putting the pieces together. Nothing |
| 7 | //! is stubbed, because the parts that would be worth stubbing are exactly the |
| 8 | //! parts that have never been run against each other before. |
| 9 | //! |
| 10 | //! What is *not* here is the handshake. There is no session and nothing is |
| 11 | //! encrypted; the packet signature says the packet was signed by the key |
| 12 | //! travelling with it and no more. That is the same bar a length-prefixed |
| 13 | //! signed envelope over TCP clears, which is the point of this round. |
| 14 | |
| 15 | use oxedyne_fe2o3_shield::srv::{ |
| 16 | cfg::ServerConfig, |
| 17 | client::Client, |
| 18 | constant, |
| 19 | context::ServerContext, |
| 20 | msg::{ |
| 21 | app::{ |
| 22 | Answer, |
| 23 | AppMsgKind, |
| 24 | }, |
| 25 | protocol::{ |
| 26 | DefaultProtocolTypes, |
| 27 | Protocol, |
| 28 | ProtocolMode, |
| 29 | }, |
| 30 | syntax as srv_syntax, |
| 31 | }, |
| 32 | schemes::WireSchemesInput, |
| 33 | server::Server, |
| 34 | }; |
| 35 | |
| 36 | use oxedyne_fe2o3_core::{ |
| 37 | prelude::*, |
| 38 | alt::Alt, |
| 39 | path::NormalPath, |
| 40 | rand::RanDef, |
| 41 | }; |
| 42 | use oxedyne_fe2o3_crypto::{ |
| 43 | enc::EncryptionScheme, |
| 44 | sign::SignatureScheme, |
| 45 | }; |
| 46 | use oxedyne_fe2o3_hash::{ |
| 47 | csum::ChecksumScheme, |
| 48 | hash::HashScheme, |
| 49 | }; |
| 50 | use oxedyne_fe2o3_net::id; |
| 51 | use oxedyne_fe2o3_o3db_sync::O3db; |
| 52 | use oxedyne_fe2o3_syntax::SyntaxRef; |
| 53 | |
| 54 | use std::{ |
| 55 | net::SocketAddr, |
| 56 | path::Path, |
| 57 | time::Duration, |
| 58 | }; |
| 59 | |
| 60 | |
| 61 | /// Length of the per-peer proof-of-work challenge code. |
| 62 | const CODE_LEN: usize = 8; |
| 63 | |
| 64 | /// Bytes per chunk. Small, so that a payload of a few kilobytes is genuinely |
| 65 | /// carried by several packets and reassembled rather than fitting in one. |
| 66 | const CHUNK_BYTES: u64 = 400; |
| 67 | |
| 68 | /// Proof-of-work difficulty, fixed: minimum and maximum are the same, so the |
| 69 | /// difficulty this peer demands does not move with the request rate. |
| 70 | /// |
| 71 | /// A difficulty that moves is only usable once a peer can be *told* the new |
| 72 | /// one, and telling it is the handshake response that is still deferred. Until |
| 73 | /// then, a fixed difficulty is the honest setting, and this test uses one |
| 74 | /// rather than pretending otherwise. |
| 75 | const POW_ZBITS: u16 = 2; |
| 76 | |
| 77 | /// How long a test waits for an answer before deciding none is coming. |
| 78 | const PATIENCE: Duration = Duration::from_secs(20); |
| 79 | |
| 80 | /// How long a test waits when it expects nothing. |
| 81 | const QUIET: Duration = Duration::from_millis(500); |
| 82 | |
| 83 | type Types = DefaultProtocolTypes<{ id::MID_LEN }, { id::SID_LEN }, { id::UID_LEN }>; |
| 84 | |
| 85 | type Proto = Protocol< |
| 86 | CODE_LEN, |
| 87 | { id::MID_LEN }, |
| 88 | { id::SID_LEN }, |
| 89 | { id::UID_LEN }, |
| 90 | Types, |
| 91 | >; |
| 92 | |
| 93 | /// The database type the server is parameterised by. No database is opened |
| 94 | /// here: the server needs the type to exist, not an instance of it. |
| 95 | type Db = O3db< |
| 96 | { id::UID_LEN }, |
| 97 | id::Uid, |
| 98 | EncryptionScheme, |
| 99 | HashScheme, |
| 100 | HashScheme, |
| 101 | ChecksumScheme, |
| 102 | >; |
| 103 | |
| 104 | type Srv = Server< |
| 105 | CODE_LEN, |
| 106 | { id::MID_LEN }, |
| 107 | { id::SID_LEN }, |
| 108 | { id::UID_LEN }, |
| 109 | Types, |
| 110 | EncryptionScheme, |
| 111 | HashScheme, |
| 112 | Db, |
| 113 | >; |
| 114 | |
| 115 | type Cli = Client< |
| 116 | CODE_LEN, |
| 117 | { id::MID_LEN }, |
| 118 | { id::SID_LEN }, |
| 119 | { id::UID_LEN }, |
| 120 | Types, |
| 121 | >; |
| 122 | |
| 123 | /// A server configuration that binds loopback on whatever port is free. |
| 124 | fn loopback_config() -> ServerConfig { |
| 125 | ServerConfig { |
| 126 | server_address: fmt!("127.0.0.1"), |
| 127 | server_port_udp: 0, |
| 128 | wire_chunk_bytes: CHUNK_BYTES, |
| 129 | wire_chunk_threshold: CHUNK_BYTES, |
| 130 | server_pow_zbits_min: POW_ZBITS, |
| 131 | server_pow_zbits_max: POW_ZBITS, |
| 132 | ..Default::default() |
| 133 | } |
| 134 | } |
| 135 | |
| 136 | /// A protocol instance with its own signing key and identifier. |
| 137 | /// |
| 138 | /// Every peer signs with a key of its own, and carries the public half in each |
| 139 | /// packet: without a session there is nowhere else for the other side to have |
| 140 | /// got it from. |
| 141 | fn protocol(cfg: &ServerConfig, uid: u128) -> Outcome<Proto> { |
| 142 | Protocol::new( |
| 143 | cfg, |
| 144 | WireSchemesInput { |
| 145 | enc: Alt::Specific(None::<EncryptionScheme>), |
| 146 | csum: Alt::Specific(None::<ChecksumScheme>), |
| 147 | powh: Alt::Specific(ServerConfig::default_packet_pow_hash_scheme()), |
| 148 | sign: Alt::Specific(Some(SignatureScheme::new_ed25519())), |
| 149 | hsenc: Alt::Specific(None::<EncryptionScheme>), |
| 150 | chnk: Some(ServerConfig::new_chunk_cfg( |
| 151 | CHUNK_BYTES as usize, |
| 152 | CHUNK_BYTES as usize, |
| 153 | false, // Do not wrap each chunk in a daticle. |
| 154 | false, // Do not pad the last chunk. |
| 155 | )), |
| 156 | }, |
| 157 | [0u8; CODE_LEN], |
| 158 | id::Mid::default(), |
| 159 | id::Sid::default(), |
| 160 | id::Uid::new(uid), |
| 161 | ProtocolMode::Test, |
| 162 | ) |
| 163 | } |
| 164 | |
| 165 | /// Build a server on loopback, hand back the address it landed on, and leave it |
| 166 | /// running with the given handler. |
| 167 | async fn serve<H, F>(handler: H) -> Outcome<SocketAddr> |
| 168 | where |
| 169 | H: Fn(Vec<u8>, SocketAddr) -> F + Send + Sync + 'static, |
| 170 | F: std::future::Future<Output = Outcome<Answer>> + Send, |
| 171 | { |
| 172 | let cfg = loopback_config(); |
| 173 | let proto = res!(protocol(&cfg, 1)); |
| 174 | let root = Path::new(".").normalise().absolute(); |
| 175 | let context = ServerContext::new(cfg, root, None::<(Db, id::Uid)>, proto); |
| 176 | let syntax = res!(srv_syntax::base_msg()); |
| 177 | let (mut server, _cmd_chan): (Srv, _) = Server::new(context, syntax); |
| 178 | let sock = res!(server.bind().await); |
| 179 | let addr = res!(sock.local_addr(), IO, Network); |
| 180 | tokio::spawn(async move { |
| 181 | if let Err(e) = server.run(sock, handler).await { |
| 182 | error!(e); |
| 183 | } |
| 184 | }); |
| 185 | Ok(addr) |
| 186 | } |
| 187 | |
| 188 | /// A client on loopback, on whatever port is free. |
| 189 | async fn dial(uid: u128) -> Outcome<(Cli, SyntaxRef)> { |
| 190 | let cfg = loopback_config(); |
| 191 | let proto = res!(protocol(&cfg, uid)); |
| 192 | let syntax = res!(srv_syntax::base_msg()); |
| 193 | let addr = res!("127.0.0.1:0".parse::<SocketAddr>(), Test); |
| 194 | let client = res!(Client::bind(addr, proto, syntax.clone()).await); |
| 195 | Ok((client, syntax)) |
| 196 | } |
| 197 | |
| 198 | /// A handler that answers with the payload it was given, reversed. |
| 199 | /// |
| 200 | /// Reversed rather than echoed, so that an answer which is somehow the |
| 201 | /// question coming back off a wire rather than out of the handler cannot pass. |
| 202 | async fn reverse(payload: Vec<u8>, _src_addr: SocketAddr) -> Outcome<Answer> { |
| 203 | Ok(Answer::Reply(payload.into_iter().rev().collect())) |
| 204 | } |
| 205 | |
| 206 | /// A payload of `n` bytes that is not the same at both ends and not the same |
| 207 | /// as any other length's. |
| 208 | fn payload(n: usize, seed: u8) -> Vec<u8> { |
| 209 | (0..n).map(|i| (i as u8).wrapping_mul(31).wrapping_add(seed)).collect() |
| 210 | } |
| 211 | |
| 212 | /// A payload larger than one chunk goes out in pieces and arrives whole, and |
| 213 | /// so does the answer to it. |
| 214 | #[tokio::test] |
| 215 | async fn test_payload_larger_than_one_chunk_00() -> Outcome<()> { |
| 216 | log_set_level!("warn"); |
| 217 | let addr = res!(serve(reverse).await); |
| 218 | let (client, _) = res!(dial(2).await); |
| 219 | |
| 220 | let sent = payload(2_000, 7); |
| 221 | assert!(sent.len() > CHUNK_BYTES as usize * 4, |
| 222 | "the point of this test is a payload of several chunks"); |
| 223 | let heard = res!(client.ask(addr, sent.clone(), PATIENCE).await); |
| 224 | |
| 225 | let mut want = sent.clone(); |
| 226 | want.reverse(); |
| 227 | assert_eq!(heard.len(), want.len(), |
| 228 | "the answer came back a different length from the question"); |
| 229 | assert_eq!(heard, want, "the answer is not what the handler made of the question"); |
| 230 | |
| 231 | // And a payload that fits in one chunk still works, so that what is being |
| 232 | // proved above is the assembly and not merely that anything at all moves. |
| 233 | let small = payload(16, 9); |
| 234 | let heard = res!(client.ask(addr, small.clone(), PATIENCE).await); |
| 235 | let mut want = small.clone(); |
| 236 | want.reverse(); |
| 237 | assert_eq!(heard, want); |
| 238 | Ok(()) |
| 239 | } |
| 240 | |
| 241 | /// Two exchanges in the air at once, and each answer finds the question that |
| 242 | /// asked it. |
| 243 | /// |
| 244 | /// The correlation is the message identifier in the packet header, and this is |
| 245 | /// what makes that load-bearing: an answer arriving under an identifier |
| 246 | /// nobody asked under is not an answer, however well-formed it is and however |
| 247 | /// much the waiting peer would like one. |
| 248 | #[tokio::test] |
| 249 | async fn test_reply_correlation_00() -> Outcome<()> { |
| 250 | log_set_level!("warn"); |
| 251 | let addr = res!(serve(reverse).await); |
| 252 | let (one, _) = res!(dial(3).await); |
| 253 | let (two, _) = res!(dial(4).await); |
| 254 | |
| 255 | let first = payload(700, 11); // Two chunks. |
| 256 | let second = payload(700, 23); // Two chunks, and different bytes. |
| 257 | assert_ne!(first, second); |
| 258 | |
| 259 | let (heard_one, heard_two) = tokio::join!( |
| 260 | one.ask(addr, first.clone(), PATIENCE), |
| 261 | two.ask(addr, second.clone(), PATIENCE), |
| 262 | ); |
| 263 | let heard_one = res!(heard_one); |
| 264 | let heard_two = res!(heard_two); |
| 265 | |
| 266 | let mut want_one = first.clone(); |
| 267 | want_one.reverse(); |
| 268 | let mut want_two = second.clone(); |
| 269 | want_two.reverse(); |
| 270 | assert_eq!(heard_one, want_one, "the first peer was given the wrong answer"); |
| 271 | assert_eq!(heard_two, want_two, "the second peer was given the wrong answer"); |
| 272 | |
| 273 | // An answer that arrives under an identifier this peer did not ask under is |
| 274 | // dropped rather than taken. The question below really is answered -- the |
| 275 | // server is alive and has just answered two others -- and the answer is |
| 276 | // waited for under a fresh identifier that nothing was ever sent with, so |
| 277 | // the only thing that can end this wait is the clock. |
| 278 | let asked = res!(one.tell(addr, payload(32, 31)).await); |
| 279 | let never = <id::Mid as RanDef>::randef(); |
| 280 | assert_ne!(never, asked, "the fresh identifier collided with a real one"); |
| 281 | match one.hear(&never, QUIET).await { |
| 282 | Ok(bytes) => return Err(err!( |
| 283 | "An answer of {} bytes was taken as the answer to a question that was \ |
| 284 | never asked.", bytes.len(); Test, Unexpected)), |
| 285 | Err(_) => (), |
| 286 | } |
| 287 | Ok(()) |
| 288 | } |
| 289 | |
| 290 | /// Rubbish sent at the port is dropped, and the peer goes on working. |
| 291 | /// |
| 292 | /// The loop is what is being tested, not the validator. A validator that |
| 293 | /// rejects a bad packet and a loop that dies doing it come to the same thing |
| 294 | /// from the outside, and the second is what an attacker would be aiming for. |
| 295 | #[tokio::test] |
| 296 | async fn test_garbage_is_dropped_00() -> Outcome<()> { |
| 297 | log_set_level!("warn"); |
| 298 | let addr = res!(serve(reverse).await); |
| 299 | |
| 300 | // Every shape of rubbish that can reach a UDP port: nothing, too little to |
| 301 | // be a header, a plausible length of noise, and a header whose chunk length |
| 302 | // claims more bytes than the packet holds. |
| 303 | let mut header_lie = payload(64, 3); |
| 304 | header_lie[0] = 0x04; // Message type 1024, an application request. |
| 305 | header_lie[1] = 0x00; |
| 306 | for i in 33..37 { |
| 307 | header_lie[i] = 0xff; // A chunk size of 65,535 in a 64-byte packet. |
| 308 | } |
| 309 | let rubbish: Vec<Vec<u8>> = vec![ |
| 310 | Vec::new(), |
| 311 | vec![0x00], |
| 312 | payload(8, 1), |
| 313 | payload(200, 2), |
| 314 | header_lie, |
| 315 | vec![0xff; constant::UDP_BUFFER_SIZE], |
| 316 | ]; |
| 317 | let noise = res!(tokio::net::UdpSocket::bind("127.0.0.1:0").await, IO, Network); |
| 318 | for r in &rubbish { |
| 319 | res!(noise.send_to(r, addr).await, IO, Network); |
| 320 | } |
| 321 | |
| 322 | // And the packet that is not rubbish at all except in one byte. This one is |
| 323 | // built by the real builder -- real header, real proof of work, real |
| 324 | // signature -- and then a single byte of the payload is flipped, which the |
| 325 | // proof of work does not cover and the signature does. The control below it |
| 326 | // is the same packet unflipped: without that, a silence here would only |
| 327 | // show that something went wrong somewhere. |
| 328 | let noise_addr = res!(noise.local_addr(), IO, Network); |
| 329 | let cfg = loopback_config(); |
| 330 | let forger = res!(protocol(&cfg, 8)); |
| 331 | let syntax = res!(srv_syntax::base_msg()); |
| 332 | let honest = res!(forger.build_app( |
| 333 | syntax.clone(), |
| 334 | AppMsgKind::Request, |
| 335 | <id::Mid as RanDef>::randef(), |
| 336 | payload(64, 17), |
| 337 | noise_addr.ip(), |
| 338 | addr.ip(), |
| 339 | )); |
| 340 | assert_eq!(honest.len(), 1, "a 64-byte payload should fit in one packet"); |
| 341 | let mut tampered = honest[0].clone(); |
| 342 | let midpoint = tampered.len() / 2; |
| 343 | tampered[midpoint] ^= 0x01; |
| 344 | res!(noise.send_to(&tampered, addr).await, IO, Network); |
| 345 | let mut heard = [0u8; constant::UDP_BUFFER_SIZE]; |
| 346 | match tokio::time::timeout(QUIET, noise.recv_from(&mut heard)).await { |
| 347 | Err(_) => (), // Nothing came back, which is the point. |
| 348 | Ok(_) => return Err(err!( |
| 349 | "A packet whose payload was altered after it was signed was answered."; |
| 350 | Test, Unexpected)), |
| 351 | } |
| 352 | |
| 353 | // The control: the same packet, unaltered, is answered. |
| 354 | let honest = res!(forger.build_app( |
| 355 | syntax, |
| 356 | AppMsgKind::Request, |
| 357 | <id::Mid as RanDef>::randef(), |
| 358 | payload(64, 17), |
| 359 | noise_addr.ip(), |
| 360 | addr.ip(), |
| 361 | )); |
| 362 | res!(noise.send_to(&honest[0], addr).await, IO, Network); |
| 363 | match tokio::time::timeout(PATIENCE, noise.recv_from(&mut heard)).await { |
| 364 | Ok(Ok(_)) => (), |
| 365 | _ => return Err(err!( |
| 366 | "The same packet unaltered was not answered, so the silence above \ |
| 367 | proves nothing."; Test, Unexpected)), |
| 368 | } |
| 369 | |
| 370 | // The peer is still there, and still answers. |
| 371 | let (client, _) = res!(dial(5).await); |
| 372 | let sent = payload(500, 13); |
| 373 | let heard = res!(client.ask(addr, sent.clone(), PATIENCE).await); |
| 374 | let mut want = sent.clone(); |
| 375 | want.reverse(); |
| 376 | assert_eq!(heard, want, "the peer stopped answering after being sent rubbish"); |
| 377 | Ok(()) |
| 378 | } |
| 379 | |
| 380 | /// A server binds where its configuration says, and loopback is a thing the |
| 381 | /// configuration can say. |
| 382 | /// |
| 383 | /// It could not before: the address was taken from the machine's network |
| 384 | /// interface whatever `server_address` held, which made two peers on one |
| 385 | /// machine -- a test, a development box -- impossible to arrange. |
| 386 | #[tokio::test] |
| 387 | async fn test_loopback_bind_00() -> Outcome<()> { |
| 388 | log_set_level!("warn"); |
| 389 | let cfg = loopback_config(); |
| 390 | let proto = res!(protocol(&cfg, 6)); |
| 391 | let root = Path::new(".").normalise().absolute(); |
| 392 | let context = ServerContext::new(cfg, root, None::<(Db, id::Uid)>, proto); |
| 393 | let syntax = res!(srv_syntax::base_msg()); |
| 394 | let (server, _cmd_chan): (Srv, _) = Server::new(context, syntax); |
| 395 | let sock = res!(server.bind().await); |
| 396 | let addr = res!(sock.local_addr(), IO, Network); |
| 397 | assert!(addr.ip().is_loopback(), |
| 398 | "a server told to bind 127.0.0.1 bound {} instead", addr.ip()); |
| 399 | assert!(addr.port() != 0, "a port of zero should have become a real one"); |
| 400 | |
| 401 | // And a setting that is neither an address nor the word for the machine's |
| 402 | // own is refused, rather than quietly becoming something else. |
| 403 | let mut bad = loopback_config(); |
| 404 | bad.server_address = fmt!("not-an-address"); |
| 405 | assert!(bad.bind_ip().is_err(), |
| 406 | "a server_address that is not an address was accepted"); |
| 407 | |
| 408 | // The word an operator writes for "wherever this machine is on its network" |
| 409 | // still resolves, which is what every deployed peer relies on. |
| 410 | let mut here = loopback_config(); |
| 411 | here.server_address = fmt!("local"); |
| 412 | res!(here.bind_ip()); |
| 413 | Ok(()) |
| 414 | } |