| 1 | //! migrate_probe CACHE...: converts copies of schema-14 caches, reports what they hold |
| 2 | //! and, where nothing blocks the queue, publishes it to an in-memory copy of the remote. |
| 3 | use notebook::{Remote, Replica}; |
| 4 | use onestore::{CommitError, CommitIo, Stamp, Transaction}; |
| 5 | use std::io; |
| 6 | |
| 7 | struct Memory(Vec<u8>); |
| 8 | |
| 9 | impl CommitIo for Memory { |
| 10 | fn read_at(&mut self, offset: u64, output: &mut [u8]) -> io::Result<usize> { |
| 11 | let offset = offset as usize; |
| 12 | let size = output.len().min(self.0.len().saturating_sub(offset)); |
| 13 | output[..size].copy_from_slice(&self.0[offset..offset + size]); |
| 14 | Ok(size) |
| 15 | } |
| 16 | fn write_at(&mut self, offset: u64, bytes: &[u8]) -> io::Result<usize> { |
| 17 | let offset = offset as usize; |
| 18 | self.0.resize(self.0.len().max(offset + bytes.len()), 0); |
| 19 | self.0[offset..offset + bytes.len()].copy_from_slice(bytes); |
| 20 | Ok(bytes.len()) |
| 21 | } |
| 22 | fn flush(&mut self) -> io::Result<()> { |
| 23 | Ok(()) |
| 24 | } |
| 25 | } |
| 26 | |
| 27 | impl Remote for Memory { |
| 28 | fn read(&mut self) -> io::Result<Vec<u8>> { |
| 29 | Ok(self.0.clone()) |
| 30 | } |
| 31 | fn stamp(&mut self) -> io::Result<Stamp> { |
| 32 | Stamp::of(&self.0).map_err(io::Error::other) |
| 33 | } |
| 34 | fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> { |
| 35 | transaction.commit(self) |
| 36 | } |
| 37 | fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> { |
| 38 | onestore::confirm(self, base) |
| 39 | } |
| 40 | } |
| 41 | |
| 42 | fn pages(replica: &Replica) -> Result<Vec<onestore::page::Page>, notebook::Error> { |
| 43 | replica |
| 44 | .pages()? |
| 45 | .into_iter() |
| 46 | .map(|(space, ..)| replica.page(space)) |
| 47 | .collect() |
| 48 | } |
| 49 | |
| 50 | fn main() -> Result<(), Box<dyn std::error::Error>> { |
| 51 | for path in std::env::args().skip(1) { |
| 52 | let copy = std::env::temp_dir().join(format!( |
| 53 | "migrate-probe-{}-{}", |
| 54 | std::process::id(), |
| 55 | std::path::Path::new(&path) |
| 56 | .file_name() |
| 57 | .unwrap() |
| 58 | .to_string_lossy() |
| 59 | )); |
| 60 | std::fs::copy(&path, &copy)?; |
| 61 | // A cache open or crashed in WAL mode keeps committed pages in its `-wal` file. |
| 62 | let wal = format!("{path}-wal"); |
| 63 | if std::path::Path::new(&wal).exists() { |
| 64 | let mut target = copy.as_os_str().to_owned(); |
| 65 | target.push("-wal"); |
| 66 | std::fs::copy(&wal, target)?; |
| 67 | } |
| 68 | let started = std::time::Instant::now(); |
| 69 | match Replica::open(&copy) { |
| 70 | Ok(replica) => { |
| 71 | let elapsed = started.elapsed(); |
| 72 | let pending = replica.pending()?; |
| 73 | let summary = replica.recovery_summary()?; |
| 74 | println!( |
| 75 | "{path}: converted in {:.1} ms; {} edits; {:?}; conflict pages {:?}", |
| 76 | elapsed.as_secs_f64() * 1e3, |
| 77 | pending.len(), |
| 78 | summary, |
| 79 | replica.conflicts()? |
| 80 | ); |
| 81 | for edit in &pending { |
| 82 | println!( |
| 83 | " {} {:?} {} ops", |
| 84 | edit.id, |
| 85 | replica.status(edit.id)?, |
| 86 | edit.edit.ops.len() |
| 87 | ); |
| 88 | } |
| 89 | for (space, title, level) in replica.pages()? { |
| 90 | println!(" page {space} {level} {title:?}"); |
| 91 | } |
| 92 | let blocked = replica.recovery_summary()?.uncertain_edits > 0; |
| 93 | if !pending.is_empty() && !blocked { |
| 94 | let local = pages(&replica)?; |
| 95 | // The remote the queue was made on, as a recovery archive records it. |
| 96 | let archive = copy.with_extension("probe-recovery"); |
| 97 | replica.export_recovery(&archive)?; |
| 98 | let mut remote = Memory(notebook::Recovery::open(&archive)?.remote_snapshot()?); |
| 99 | std::fs::remove_file(&archive)?; |
| 100 | let synced = replica.sync_once(&mut remote)?; |
| 101 | let published = { |
| 102 | let arena = onestore::Arena::default(); |
| 103 | let mut section = onestore::Section::open(&arena, remote.0.clone())?; |
| 104 | section |
| 105 | .pages()? |
| 106 | .into_iter() |
| 107 | .map(|(space, ..)| section.page(space)) |
| 108 | .collect::<Result<Vec<_>, _>>()? |
| 109 | }; |
| 110 | println!( |
| 111 | " published {:?}; remote pages equal the converted pages: {}", |
| 112 | synced.edit.map(|(id, _)| id), |
| 113 | published == local |
| 114 | ); |
| 115 | } |
| 116 | } |
| 117 | Err(error) => println!("{path}: {error}"), |
| 118 | } |
| 119 | } |
| 120 | Ok(()) |
| 121 | } |