| 1 | mod support { |
| 2 | pub mod view; |
| 3 | } |
| 4 | use support::view::{cached, view}; |
| 5 | |
| 6 | use notebook::{EditStatus, Remote, Replica}; |
| 7 | use onestore::{CommitError, CommitIo, Stamp, Transaction}; |
| 8 | use serde_json::json; |
| 9 | use std::{ |
| 10 | env, |
| 11 | fs::{self, File, OpenOptions}, |
| 12 | io::{self, Write}, |
| 13 | os::unix::fs::FileExt, |
| 14 | path::Path, |
| 15 | }; |
| 16 | |
| 17 | fn phase(name: &str) { |
| 18 | println!("{}", json!({"event":"phase", "name":name})); |
| 19 | io::stdout().flush().unwrap(); |
| 20 | if env::var("ONESTORE_RECOVERY_PAUSE").ok().as_deref() == Some(name) { |
| 21 | let mut line = String::new(); |
| 22 | assert!( |
| 23 | io::stdin().read_line(&mut line).unwrap() > 0, |
| 24 | "Controller closed a paused operation" |
| 25 | ); |
| 26 | } |
| 27 | } |
| 28 | |
| 29 | struct Disk { |
| 30 | file: File, |
| 31 | writes: usize, |
| 32 | flushes: usize, |
| 33 | } |
| 34 | |
| 35 | impl CommitIo for Disk { |
| 36 | fn read_at(&mut self, offset: u64, output: &mut [u8]) -> io::Result<usize> { |
| 37 | self.file.read_at(output, offset) |
| 38 | } |
| 39 | fn write_at(&mut self, offset: u64, data: &[u8]) -> io::Result<usize> { |
| 40 | self.writes += 1; |
| 41 | println!( |
| 42 | "{}", |
| 43 | json!({"event":"write", "number":self.writes, "offset":offset, "bytes":data.len()}) |
| 44 | ); |
| 45 | phase(&format!("write-{}-before", self.writes)); |
| 46 | let result = self.file.write_at(data, offset); |
| 47 | phase(&format!("write-{}-after", self.writes)); |
| 48 | result |
| 49 | } |
| 50 | fn flush(&mut self) -> io::Result<()> { |
| 51 | self.flushes += 1; |
| 52 | phase(&format!("flush-{}-before", self.flushes)); |
| 53 | let result = self.file.sync_all(); |
| 54 | phase(&format!("flush-{}-after", self.flushes)); |
| 55 | result |
| 56 | } |
| 57 | } |
| 58 | |
| 59 | impl Remote for Disk { |
| 60 | fn read(&mut self) -> io::Result<Vec<u8>> { |
| 61 | phase("read-before"); |
| 62 | let result = onestore::read_snapshot( |
| 63 | |offset, output| self.file.read_at(output, offset), |
| 64 | 256 * 1024 * 1024, |
| 65 | ) |
| 66 | .and_then(|snapshot| snapshot.ok_or_else(|| io::ErrorKind::WouldBlock.into())); |
| 67 | phase("read-after"); |
| 68 | result |
| 69 | } |
| 70 | fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> { |
| 71 | phase("publish-before"); |
| 72 | let result = transaction.commit(self); |
| 73 | phase("publish-after"); |
| 74 | result |
| 75 | } |
| 76 | fn stamp(&mut self) -> io::Result<Stamp> { |
| 77 | Stamp::of(&self.read()?).map_err(io::Error::other) |
| 78 | } |
| 79 | fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> { |
| 80 | phase("confirm-before"); |
| 81 | let result = onestore::confirm(self, base); |
| 82 | phase("confirm-after"); |
| 83 | result |
| 84 | } |
| 85 | } |
| 86 | |
| 87 | fn report(cache: &Replica, root: &Path) -> Result<(), Box<dyn std::error::Error>> { |
| 88 | let (_, _, local) = cached(cache)?; |
| 89 | let remote = view(&fs::read(root.join("remote.one"))?)?; |
| 90 | let (status, revision) = match cache.status(1)? { |
| 91 | Some(EditStatus::Pending) => ("pending", None), |
| 92 | Some(EditStatus::AwaitingConfirmation { revision }) => { |
| 93 | ("uncertain", Some(revision.to_string())) |
| 94 | } |
| 95 | Some(EditStatus::Published { revision }) => ("published", Some(revision.to_string())), |
| 96 | Some(EditStatus::Archived { archive }) => ("archived", Some(archive)), |
| 97 | None => ("missing", None), |
| 98 | }; |
| 99 | println!( |
| 100 | "{}", |
| 101 | json!({"event":"state", "status":status, "revision":revision, |
| 102 | "local_text":local, "remote_text":remote.text, "remote_revision":remote.revision.to_string(), |
| 103 | "pending":cache.pending()?.iter().map(|pending| match &pending.edit.ops[..] { |
| 104 | [onestore::op::Op::Page { op: onestore::op::PageOp::Text { range, with, .. }, .. }] => |
| 105 | json!({"id":pending.id,"replacement":with,"range":[range.start,range.end]}), |
| 106 | _ => json!({"id":pending.id,"edit":pending.edit}), |
| 107 | }).collect::<Vec<_>>() }) |
| 108 | ); |
| 109 | Ok(()) |
| 110 | } |
| 111 | |
| 112 | fn main() -> Result<(), Box<dyn std::error::Error>> { |
| 113 | let args: Vec<_> = env::args().skip(1).collect(); |
| 114 | let mode = args |
| 115 | .first() |
| 116 | .ok_or("Expected init|inspect|sync DIRECTORY [SOURCE]")?; |
| 117 | let root = Path::new(args.get(1).ok_or("Missing owned directory")?); |
| 118 | if mode == "init" && args.len() == 3 { |
| 119 | let source = fs::read(&args[2])?; |
| 120 | let target = view(&source)?; |
| 121 | fs::create_dir(root)?; |
| 122 | let mut remote = OpenOptions::new() |
| 123 | .write(true) |
| 124 | .create_new(true) |
| 125 | .open(root.join("remote.one"))?; |
| 126 | remote.write_all(&source)?; |
| 127 | remote.sync_all()?; |
| 128 | drop(remote); |
| 129 | let cache = Replica::create(root.join("cache.sqlite"), &source)?; |
| 130 | let at = u32::try_from(target.text.encode_utf16().count())?; |
| 131 | let edit = onestore::op::Edit { |
| 132 | at: 133_000_000_000_000_000, |
| 133 | ops: vec![onestore::op::Op::Page { |
| 134 | space: target.space, |
| 135 | op: onestore::op::PageOp::Text { |
| 136 | text: target.object, |
| 137 | range: at..at, |
| 138 | with: " [offline-recovery]".into(), |
| 139 | }, |
| 140 | }], |
| 141 | }; |
| 142 | assert_eq!(cache.apply("Fixture", edit)?, 1); |
| 143 | phase("local-after"); |
| 144 | report(&cache, root)?; |
| 145 | } else if args.len() == 2 && ["inspect", "sync"].contains(&mode.as_str()) { |
| 146 | let cache = Replica::open(root.join("cache.sqlite"))?; |
| 147 | if mode == "sync" { |
| 148 | let mut disk = Disk { |
| 149 | file: OpenOptions::new() |
| 150 | .read(true) |
| 151 | .write(true) |
| 152 | .open(root.join("remote.one"))?, |
| 153 | writes: 0, |
| 154 | flushes: 0, |
| 155 | }; |
| 156 | phase("sync-before"); |
| 157 | let result = cache.sync_once(&mut disk); |
| 158 | phase("sync-after"); |
| 159 | result?; |
| 160 | } |
| 161 | report(&cache, root)?; |
| 162 | } else { |
| 163 | return Err("Expected init|inspect|sync DIRECTORY [SOURCE]".into()); |
| 164 | } |
| 165 | Ok(()) |
| 166 | } |