oxedyne/fe2o3/fe2o3_shield/tests/sim.rs
4.5 KiB, 46 runs
created by r1870400018:5006, 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 | use oxedyne_fe2o3_shield::{ |
| 2 | //prelude::*, |
| 3 | app::{ |
| 4 | cfg::AppConfig, |
| 5 | constant, |
| 6 | server, |
| 7 | tui::AppStatus, |
| 8 | }, |
| 9 | srv::{ |
| 10 | cmd::Command, |
| 11 | context::new_db, |
| 12 | }, |
| 13 | }; |
| 14 | |
| 15 | use oxedyne_fe2o3_core::{ |
| 16 | prelude::*, |
| 17 | channels::Recv, |
| 18 | log::{ |
| 19 | console::{ |
| 20 | switch_to_logger_console, |
| 21 | MultiStreamLoggerConsole, |
| 22 | StdoutLoggerConsole, |
| 23 | }, |
| 24 | }, |
| 25 | }; |
| 26 | use oxedyne_fe2o3_hash::kdf::KeyDerivationScheme; |
| 27 | use oxedyne_fe2o3_iop_hash::kdf::KeyDeriver; |
| 28 | |
| 29 | use std::{ |
| 30 | fs, |
| 31 | path::Path, |
| 32 | time::{ |
| 33 | Duration, |
| 34 | Instant, |
| 35 | }, |
| 36 | thread, |
| 37 | }; |
| 38 | |
| 39 | |
| 40 | pub fn test_sim(_filter: &'static str) -> Outcome<()> { |
| 41 | let rt = res!(tokio::runtime::Runtime::new()); |
| 42 | rt.block_on(run_test_sim(_filter)) |
| 43 | } |
| 44 | |
| 45 | pub async fn run_test_sim(_filter: &'static str) -> Outcome<()> { |
| 46 | |
| 47 | test!("Reconfiguring log to stream to multiple channels..."); |
| 48 | res!(switch_to_logger_console::<MultiStreamLoggerConsole<_>>()); |
| 49 | let mut log_cfg = log_get_config!(); |
| 50 | log_cfg.file = None; |
| 51 | log_set_config!(log_cfg); |
| 52 | |
| 53 | const NUM_PEERS: usize = 1; |
| 54 | |
| 55 | let mut server_handles = Vec::new(); |
| 56 | |
| 57 | for peer in 1..=NUM_PEERS { |
| 58 | let id = fmt!("{:03}", peer); |
| 59 | msg!("Starting peer: {}", id); |
| 60 | let peer_dir = fmt!("sims/{}", id); |
| 61 | |
| 62 | if res!(fs::exists(&peer_dir)) { |
| 63 | res!(fs::remove_dir_all(&peer_dir)); |
| 64 | } |
| 65 | res!(fs::create_dir_all(&peer_dir)); |
| 66 | |
| 67 | let mut peer_cfg = res!(AppConfig::new()); |
| 68 | peer_cfg.app_root = peer_dir.clone(); |
| 69 | |
| 70 | let app_root = Path::new(&peer_cfg.app_root); |
| 71 | let db_root = app_root.join(constant::DB_DIR); |
| 72 | |
| 73 | let mut db_kdf = res!(KeyDerivationScheme::from_str(&peer_cfg.kdf_name)); |
| 74 | res!(db_kdf.derive(id.as_bytes())); // Use the id as the db encryption passphrase. |
| 75 | let db_enc_key = res!(db_kdf.get_hash()).to_vec(); |
| 76 | |
| 77 | match res!(server::start_server( |
| 78 | &peer_cfg, |
| 79 | &AppStatus::default(), |
| 80 | res!(new_db(&db_root, &db_enc_key)), |
| 81 | None, |
| 82 | Some(id.clone()), |
| 83 | ).await) { |
| 84 | (_eval, Some((cmd_chan, handle))) => { |
| 85 | server_handles.push((id, cmd_chan, handle)); |
| 86 | } |
| 87 | (_eval, opt) => { |
| 88 | return Err(err!( |
| 89 | "Unexpected server start response: opt = {:?}", opt; |
| 90 | Test, Unexpected)); |
| 91 | } |
| 92 | } |
| 93 | } |
| 94 | |
| 95 | thread::sleep(Duration::from_secs(5)); |
| 96 | |
| 97 | let id = fmt!("001"); |
| 98 | match log_get_streams!(Duration::from_secs(2)) { |
| 99 | Ok(Some(streams_map)) => { |
| 100 | let unlocked_map = lock_read!(streams_map); |
| 101 | if let Some(log_stream_chan) = unlocked_map.get(&id) { |
| 102 | // Print out log stream for 5s. |
| 103 | let duration = Duration::from_secs(5); |
| 104 | let send_finish = Duration::from_secs(4); |
| 105 | let start_time = Instant::now(); |
| 106 | let mut finish_sent = false; |
| 107 | |
| 108 | while start_time.elapsed() < duration { |
| 109 | match log_stream_chan.try_recv() { |
| 110 | Recv::Empty => {} |
| 111 | Recv::Result(Ok(line)) => { |
| 112 | println!("{}", line); |
| 113 | } |
| 114 | Recv::Result(Err(e)) => return Err(err!(e, |
| 115 | "While reading command channel."; Channel, Read)), |
| 116 | } |
| 117 | if !finish_sent && start_time.elapsed() > send_finish { |
| 118 | for (_, cmd_chan, _) in &server_handles { |
| 119 | res!(cmd_chan.send(Command::Finish)); |
| 120 | } |
| 121 | finish_sent = true; |
| 122 | } |
| 123 | } |
| 124 | } |
| 125 | } |
| 126 | Ok(None) => { |
| 127 | return Err(err!("Could not acquire log streams map."; Channel, Missing)); |
| 128 | } |
| 129 | Err(e) => { |
| 130 | return Err(err!(e, "While requesting log streams map."; Channel)); |
| 131 | } |
| 132 | } |
| 133 | |
| 134 | res!(switch_to_logger_console::<StdoutLoggerConsole<_>>()); |
| 135 | |
| 136 | thread::sleep(Duration::from_secs(1)); |
| 137 | |
| 138 | // Wait for all servers to finish |
| 139 | for (id, _, handle) in server_handles { |
| 140 | match handle.await { |
| 141 | Ok(server_result) => match server_result { |
| 142 | Ok(()) => info!("Server {} stopped gracefully.", id), |
| 143 | Err(e) => error!(err!(e, |
| 144 | "While running server {} task.", id; IO, Thread)), |
| 145 | }, |
| 146 | Err(join_err) => return Err(err!(join_err, |
| 147 | "While awaiting server {} task completion.", id; Async)), |
| 148 | } |
| 149 | } |
| 150 | |
| 151 | Ok(()) |
| 152 | } |