| 719 | } |
| 720 | |
| 721 | async fn atomically<F, Fut>(file_path: Option<PathBuf>, f: F) -> Result<()> |
| 722 | where |
| 723 | F: FnOnce(AtomicWriter) -> Fut, |
| 724 | Fut: Future<Output = Result<AtomicWriter>>, |
| 725 | { |
| 726 | match file_path { |
| 727 | Some(file_path) => { |
| 728 | let dir = file_path.parent().expect("file not in a directory").to_owned(); |
| 729 | fs::create_dir_all(&dir).await?; |
| 730 | let (tmp_file, tmp_out) = spawn_blocking({ |
| 731 | let dir = dir.clone(); |
| 732 | move || { |
| 733 | let tmp = NamedTempFile::new_in(dir)?; |
| 734 | let out = tmp.reopen()?; |
| 735 | Ok::<_, io::Error>((tmp, out)) |
| 736 | } |
| 737 | }) |
| 738 | .await |
| 739 | .unwrap()?; |
| 740 | |
| 741 | let mut file = AtomicWriter::File(fs::File::from_std(tmp_out)); |
| 742 | file = f(file).await?; |
| 743 | file.sync_all().await?; |
| 744 | |
| 745 | spawn_blocking(|| tmp_file.persist(file_path)) |
| 746 | .await |
| 747 | .unwrap() |
| 748 | .map_err(|e| e.error)?; |
| 749 | fs::File::open(dir).await?.sync_all().await?; |
| 750 | } |
| 751 | |
| 752 | None => { |
| 753 | f(AtomicWriter::Null(tokio::io::sink())).await?; |
| 754 | } |
| 755 | } |
| 756 | |
| 757 | Ok(()) |
| 758 | } |
| 759 | |
| 760 | fn serialize_snapshot(w: &mut impl SatsWriter, value: &Snapshot) -> Result<()> { |
| 761 | bsatn::to_writer(w, value).map_err(|cause| SnapshotError::Serialize { |