1mod support {
2 pub mod view;
3}
4use support::view::{cached, view};
5
6use notebook::{EditStatus, Remote, Replica};
7use onestore::{CommitError, CommitIo, Stamp, Transaction};
8use serde_json::json;
9use std::{
10 env,
11 fs::{self, File, OpenOptions},
12 io::{self, Write},
13 os::unix::fs::FileExt,
14 path::Path,
15};
16
17fn 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
29struct Disk {
30 file: File,
31 writes: usize,
32 flushes: usize,
33}
34
35impl 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
59impl 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
87fn 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
112fn 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}