Receive model data from master and write to the cache directory.
(
stream: &mut TcpStream,
cache_dir: &Path,
layers: &[String],
on_progress: Option<&SetupProgressFn>,
)
| 938 | |
| 939 | /// Receive model data from master and write to the cache directory. |
| 940 | async fn receive_model_data( |
| 941 | stream: &mut TcpStream, |
| 942 | cache_dir: &Path, |
| 943 | layers: &[String], |
| 944 | on_progress: Option<&SetupProgressFn>, |
| 945 | ) -> Result<()> { |
| 946 | let overall_start = Instant::now(); |
| 947 | let mut overall_bytes: u64 = 0; |
| 948 | let mut current_file: Option<(String, std::fs::File, Instant, u64)> = None; |
| 949 | |
| 950 | let layer_range = if layers.is_empty() { |
| 951 | "(none)".to_string() |
| 952 | } else { |
| 953 | format!( |
| 954 | "{} — {} ({} layers)", |
| 955 | layers.first().unwrap(), |
| 956 | layers.last().unwrap(), |
| 957 | layers.len() |
| 958 | ) |
| 959 | }; |
| 960 | |
| 961 | log::info!("receiving model data [{}] ...", layer_range); |
| 962 | |
| 963 | let mut read_buf = Vec::with_capacity(MODEL_DATA_CHUNK_SIZE + 1024); |
| 964 | loop { |
| 965 | let (_, msg) = Message::from_reader_buf(stream, &mut read_buf).await?; |
| 966 | |
| 967 | match msg { |
| 968 | Message::ModelDataChunk { |
| 969 | filename, |
| 970 | offset, |
| 971 | total_size, |
| 972 | compressed, |
| 973 | checksum, |
| 974 | data, |
| 975 | } => { |
| 976 | // Verify CRC32 checksum before writing |
| 977 | let actual_crc = crc32fast::hash(&data); |
| 978 | if actual_crc != checksum { |
| 979 | return Err(anyhow!( |
| 980 | "checksum mismatch for {} at offset {}: expected {:#x}, got {:#x}", |
| 981 | filename, offset, checksum, actual_crc |
| 982 | )); |
| 983 | } |
| 984 | |
| 985 | // Decompress if compressed |
| 986 | let data = if compressed { |
| 987 | zstd::decode_all(data.as_slice()) |
| 988 | .map_err(|e| anyhow!("zstd decompress failed for {} at offset {}: {}", filename, offset, e))? |
| 989 | } else { |
| 990 | data |
| 991 | }; |
| 992 | // Open new file if needed |
| 993 | let file = if let Some((ref name, ref mut file, _, _)) = current_file { |
| 994 | if name == &filename { |
| 995 | file |
| 996 | } else { |
| 997 | // Close previous file, log stats |
no test coverage detected