1#[path = "../../onestore/examples/support/concurrent.rs"]
2mod concurrent;
3
4use notebook::smb::{Client, Credentials};
5use notebook::{EditStatus, Error, Remote, Replica, SmbRemote};
6use onestore::{CommitError, CommitState, ExGuid, RevisionIndex, Stamp, Store, Transaction};
7use serde_json::json;
8use std::{
9 env,
10 io::{self, Write},
11 path::{Path, PathBuf},
12 sync::Arc,
13 thread,
14 time::{Duration, Instant, SystemTime, UNIX_EPOCH},
15};
16
17mod support {
18 pub mod view;
19}
20use concurrent::{DocumentView, document_view};
21use support::view::{cached, view};
22
23impl DocumentView {
24 fn changes(&self, before: &Self) -> Self {
25 let difference =
26 |old: &serde_json::Map<String, serde_json::Value>,
27 new: &serde_json::Map<String, serde_json::Value>| {
28 old.keys()
29 .chain(new.keys())
30 .filter(|id| old.get(*id) != new.get(*id))
31 .map(|id| (id.clone(), new.get(id).cloned().unwrap_or_default()))
32 .collect()
33 };
34 Self {
35 texts: difference(&before.texts, &self.texts),
36 graph: difference(&before.graph, &self.graph),
37 }
38 }
39}
40
41fn now() -> u128 {
42 SystemTime::now()
43 .duration_since(UNIX_EPOCH)
44 .unwrap()
45 .as_micros()
46}
47
48#[derive(Clone)]
49enum Pause {
50 Outage(PathBuf),
51 FormatReply(PathBuf),
52}
53
54struct Traced<R> {
55 remote: R,
56 before: Option<(String, Option<DocumentView>)>,
57 /// The last image read, which publications apply to.
58 read: Vec<u8>,
59 pause: Option<Pause>,
60 documents: bool,
61}
62
63impl<R: Remote> Remote for Traced<R> {
64 fn read(&mut self) -> io::Result<Vec<u8>> {
65 let started = now();
66 let bytes = self.remote.read()?;
67 let observed = view(&bytes).map_err(|error| io::Error::other(error.to_string()))?;
68 let documents = if self.documents {
69 Some(document_view(&bytes).map_err(io::Error::other)?)
70 } else {
71 None
72 };
73 println!(
74 "{}",
75 json!({"event":"read", "started_us":started, "finished_us":now(), "text":observed.text, "documents":documents.as_ref().map(|view| &view.texts), "document_graph":documents.as_ref().map(|view| &view.graph)})
76 );
77 self.before = Some((observed.text, documents));
78 self.read.clone_from(&bytes);
79 Ok(bytes)
80 }
81 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
82 // Publication reads nothing while the stamp holds; observe an unseen base once.
83 if Stamp::of(&self.read).ok().as_ref() != Some(transaction.base()) {
84 self.read().map_err(|error| CommitError {
85 state: CommitState::NotCommitted,
86 error,
87 })?;
88 }
89 let mut image = self.read.clone();
90 transaction.apply(&mut image).map_err(|error| CommitError {
91 state: CommitState::NotCommitted,
92 error: io::Error::other(error.to_string()),
93 })?;
94 let after = view(&image).map_err(|error| CommitError {
95 state: CommitState::NotCommitted,
96 error: io::Error::other(error.to_string()),
97 })?;
98 let documents = if self.documents {
99 Some(document_view(&image).map_err(|error| CommitError {
100 state: CommitState::NotCommitted,
101 error: io::Error::other(error),
102 })?)
103 } else {
104 None
105 };
106 let changes = if let Some(after) = &documents {
107 let before = self
108 .before
109 .as_ref()
110 .and_then(|(_, documents)| documents.as_ref())
111 .ok_or_else(|| CommitError {
112 state: CommitState::NotCommitted,
113 error: io::Error::other("Document publication has no observed source"),
114 })?;
115 Some(after.changes(before))
116 } else {
117 None
118 };
119 let pause = match &self.pause {
120 Some(Pause::Outage(marker)) => Some((
121 marker,
122 marker.parent().unwrap().join("offline-outage-resumed"),
123 "outage",
124 )),
125 Some(Pause::FormatReply(marker))
126 if changes.as_ref().is_some_and(|changes| {
127 changes.texts.len() == 1
128 && self
129 .before
130 .as_ref()
131 .and_then(|(_, before)| before.as_ref())
132 .is_some_and(|before| {
133 changes.texts.keys().all(|id| before.texts.contains_key(id))
134 })
135 }) =>
136 {
137 Some((marker, marker.with_extension("resume"), "format"))
138 }
139 _ => None,
140 };
141 if let Some((marker, resumed, kind)) = pause
142 && !resumed.exists()
143 {
144 std::fs::write(marker, after.revision.to_string()).map_err(|error| CommitError {
145 state: CommitState::NotCommitted,
146 error,
147 })?;
148 println!(
149 "{}",
150 json!({"event":"publication_paused", "revision":after.revision.to_string(), "kind":kind, "at_us":now()})
151 );
152 let deadline = Instant::now() + Duration::from_secs(120);
153 while !resumed.exists() {
154 if Instant::now() >= deadline {
155 return Err(CommitError {
156 state: CommitState::NotCommitted,
157 error: io::Error::new(
158 io::ErrorKind::TimedOut,
159 "Offline publication barrier timed out",
160 ),
161 });
162 }
163 thread::sleep(Duration::from_millis(10));
164 }
165 }
166 let started = now();
167 let result = self.remote.publish(transaction);
168 let finished = now();
169 println!(
170 "{}",
171 json!({"event":"remote_attempt", "started_us":started, "finished_us":finished,
172 "revision":after.revision.to_string(), "space":after.space.to_string(), "object":after.object.to_string(), "before":self.before.as_ref().map(|(text,_)|text), "after":after.text,
173 "state":format!("{:?}", result.as_ref().map_or_else(|error| error.state, |_| CommitState::Committed)),
174 "documents":documents.as_ref().map(|view| &view.texts), "document_graph":documents.as_ref().map(|view| &view.graph),
175 "document_changes":changes.as_ref().map(|view| &view.texts), "document_graph_changes":changes.as_ref().map(|view| &view.graph)})
176 );
177 if result
178 .as_ref()
179 .is_err_and(|error| error.state == CommitState::Unknown)
180 && let Some(Pause::FormatReply(marker)) = &self.pause
181 && marker.with_extension("isolate").exists()
182 {
183 std::fs::write(marker.with_extension("isolate"), serde_json::to_vec(&json!({"space":after.space.to_string(), "revision":after.revision.to_string(), "after_us":finished})).unwrap())
184 .map_err(|error| CommitError { state:CommitState::Unknown, error })?;
185 println!(
186 "{}",
187 json!({"event":"confirmation_paused", "revision":after.revision.to_string(), "at_us":now()})
188 );
189 let deadline = Instant::now() + Duration::from_secs(120);
190 while !marker.with_extension("confirmation-resume").exists() {
191 if Instant::now() >= deadline {
192 return Err(CommitError {
193 state: CommitState::Unknown,
194 error: io::Error::new(
195 io::ErrorKind::TimedOut,
196 "Offline confirmation barrier timed out",
197 ),
198 });
199 }
200 thread::sleep(Duration::from_millis(10));
201 }
202 }
203 if result.is_ok() {
204 self.read = image;
205 self.before = Some((after.text, documents));
206 }
207 result
208 }
209 fn stamp(&mut self) -> io::Result<Stamp> {
210 self.remote.stamp()
211 }
212 fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> {
213 let captured = (|| -> Result<_, Box<dyn std::error::Error>> {
214 // A confirmation follows the read of its image, except for edits that changed nothing.
215 if Stamp::of(&self.read).ok().as_ref() != Some(base) {
216 self.read = self.remote.read()?;
217 }
218 let snapshot = &self.read[..];
219 let store = Store::parse(snapshot)?;
220 let index = RevisionIndex::parse(&store)?;
221 let revisions = index
222 .spaces
223 .iter()
224 .map(|(id, space)| {
225 (
226 id.to_string(),
227 space
228 .revisions
229 .keys()
230 .map(ToString::to_string)
231 .collect::<Vec<_>>(),
232 )
233 })
234 .collect::<std::collections::BTreeMap<_, _>>();
235 let current = index
236 .spaces
237 .iter()
238 .filter_map(|(id, space)| {
239 space
240 .labels
241 .get(&(ExGuid::default(), 1))
242 .map(|revision| (id.to_string(), revision.to_string()))
243 })
244 .collect::<std::collections::BTreeMap<_, _>>();
245 let path = env::var_os("ONESTORE_OFFLINE_CONFIRM_DIR").map(|directory| {
246 PathBuf::from(directory).join(format!("{}-{}.one", std::process::id(), now()))
247 });
248 if let Some(path) = &path {
249 std::fs::OpenOptions::new()
250 .write(true)
251 .create_new(true)
252 .open(path)?
253 .write_all(snapshot)?;
254 }
255 Ok((revisions, current, path))
256 })();
257 let (revisions, current, capture) = captured.map_err(|error| CommitError {
258 state: CommitState::NotCommitted,
259 error: io::Error::other(error.to_string()),
260 })?;
261 let started = now();
262 let result = self.remote.confirm(base);
263 println!(
264 "{}",
265 json!({"event":"remote_confirm", "started_us":started, "finished_us":now(), "revisions":revisions, "current_revisions":current, "capture":capture.as_ref().and_then(|path| path.file_name()).map(|name| name.to_string_lossy()), "text":self.before.as_ref().map(|(text,_)|text),
266 "state":format!("{:?}", result.as_ref().map_or_else(|error| error.state, |_| CommitState::Committed)), "error":result.as_ref().err().map(|error|error.error.to_string())})
267 );
268 result
269 }
270}
271
272/// The edit appending `token` at `at` of a text.
273fn append(space: ExGuid, text: ExGuid, at: u32, token: &str) -> onestore::op::Edit {
274 onestore::op::Edit {
275 at: u64::try_from(now()).unwrap_or_default() * 10 + 116_444_736_000_000_000,
276 ops: vec![onestore::op::Op::Page {
277 space,
278 op: onestore::op::PageOp::Text {
279 text,
280 range: at..at,
281 with: token.to_owned(),
282 },
283 }],
284 }
285}
286
287fn main() -> Result<(), Box<dyn std::error::Error>> {
288 let args: Vec<_> = env::args().skip(1).collect();
289 if args.len() != 7
290 || !["read", "write"].contains(&args[0].as_str())
291 || !args[2].bytes().all(|b| b.is_ascii_alphanumeric())
292 {
293 return Err("Expected read|write FILE ACTOR OPERATIONS START_FILE STOP_FILE SEED".into());
294 }
295 let documents = env::var_os("ONESTORE_OFFLINE_DOCUMENTS").is_some();
296 let address = env::var("ONESTORE_SMB_LAB")?;
297 let share = env::var("ONESTORE_SMB_SHARE")?;
298 let initial = Client::connect(
299 &address,
300 &share,
301 Credentials::default(),
302 Duration::from_secs(5),
303 )?;
304 if args[0] == "read" {
305 return concurrent::run(
306 &args,
307 |path| initial.read(path, 256 * 1024 * 1024),
308 |_, _, _, _, _, _| unreachable!(),
309 );
310 }
311 let timeout: u64 = env::var("ONESTORE_CLIENT_TIMEOUT_MS")
312 .unwrap_or_else(|_| "600000".into())
313 .parse()?;
314 let deadline = Instant::now()
315 .checked_add(Duration::from_millis(timeout))
316 .ok_or("Invalid timeout")?;
317 let operations: usize = args[3].parse()?;
318 let mut seed: u64 = args[6].parse()?;
319 if operations == 0 {
320 return Err("Expected positive operations".into());
321 }
322 let source = initial.read(&args[1], 256 * 1024 * 1024)?;
323 drop(initial);
324 let cache_path = Path::new(&args[4])
325 .parent()
326 .ok_or("Missing workload directory")?
327 .join(format!("{}.sqlite", args[2]));
328 let cache = Arc::new(Replica::create(&cache_path, &source)?);
329 println!(
330 "{}",
331 json!({"event":"ready", "pid":std::process::id(), "actor":args[2], "offline":true, "document_operations":false,"document_graph":documents,"document_kinds":[]})
332 );
333 while !Path::new(&args[4]).exists() {
334 if Instant::now() >= deadline {
335 return Err("Start barrier timed out".into());
336 }
337 thread::sleep(Duration::from_millis(5));
338 }
339 let path = args[1].clone();
340 let outage = env::var_os("ONESTORE_OFFLINE_OUTAGE_DIR").map(PathBuf::from);
341 let pause = outage
342 .as_ref()
343 .map(|directory| Pause::Outage(directory.join(format!("offline-paused-{}", args[2]))))
344 .or_else(|| {
345 env::var_os("ONESTORE_OFFLINE_FORMAT_REPLY_DIR").map(|directory| {
346 Pause::FormatReply(
347 PathBuf::from(directory).join(format!("offline-paused-{}", args[2])),
348 )
349 })
350 });
351 let (fatal_tx, fatal_rx) = std::sync::mpsc::channel();
352 let worker = cache.start_sync(Duration::from_millis(50), move || {
353 let client = Client::connect(&address, &share, Credentials::default(), Duration::from_secs(5))?;
354 println!("{}", json!({"event":"transport_connected", "at_us":now()}));
355 Ok(Traced { remote: SmbRemote::new(client, &path, 256 * 1024 * 1024), before:None, read:Vec::new(), pause:pause.clone(), documents })
356 }, move |result| {
357 if let Err(error) = result {
358 println!("{}", json!({"event":"sync_error", "error":error.to_string(), "at_us":now()}));
359 if !matches!(error, Error::RemoteIo(_) | Error::Remote(_)) && !matches!(error, Error::Io(error) if error.kind() == io::ErrorKind::WouldBlock) {
360 let _ = fatal_tx.send(error.to_string());
361 }
362 }
363 })?;
364 let mut ids = Vec::new();
365 let mut tokens = Vec::new();
366 let result = (|| -> Result<(), Box<dyn std::error::Error>> {
367 let mut generated = 0;
368 let mut received = 0;
369 loop {
370 match fatal_rx.try_recv() {
371 Ok(error) => return Err(error.into()),
372 Err(std::sync::mpsc::TryRecvError::Disconnected) => {
373 return Err("Worker exited before completion".into());
374 }
375 Err(std::sync::mpsc::TryRecvError::Empty) => {}
376 }
377 while received < ids.len() {
378 match cache.status(ids[received])? {
379 Some(EditStatus::Published { revision }) => {
380 println!(
381 "{}",
382 json!({"event":"remote_receipt", "id":ids[received], "revision":revision.to_string(), "at_us":now()})
383 );
384 received += 1;
385 }
386 Some(_) => break,
387 None => return Err("Local intent disappeared".into()),
388 }
389 }
390 if Instant::now() >= deadline {
391 return Err("Offline workload timed out; cache retained".into());
392 }
393 let pending = cache.pending()?;
394 let capacity = if outage
395 .as_ref()
396 .is_some_and(|directory| !directory.join("offline-outage-down").exists())
397 {
398 1
399 } else {
400 8
401 };
402 if generated < operations && pending.len() < capacity {
403 let (space, object, text) = cached(&cache)?;
404 let at = u32::try_from(text.encode_utf16().count())?;
405 let token = format!(" [{}:{}]", args[2], generated);
406 let started = now();
407 match cache.apply("Offline document writer", append(space, object, at, &token)) {
408 Ok(id) => {
409 println!(
410 "{}",
411 json!({"event":"local_commit", "id":id, "operation":generated, "space":space.to_string(), "object":object.to_string(), "before":text, "token":token, "started_us":started, "finished_us":now()})
412 );
413 ids.push(id);
414 tokens.push(token);
415 generated += 1;
416 }
417 other => return Err(format!("Unexpected local result: {other:?}").into()),
418 }
419 }
420 if received == ids.len() && cache.pending()?.is_empty() {
421 // Concurrent appends meet at the end of the text: a merge keeps the remote's
422 // text and this writer's on a conflict page. Append the published tokens
423 // the text lacks after it again.
424 let (space, object, text) = cached(&cache)?;
425 let missing: Vec<String> = tokens
426 .iter()
427 .filter(|token| !text.contains(token.as_str()))
428 .cloned()
429 .collect();
430 let mut at = u32::try_from(text.encode_utf16().count())?;
431 for token in &missing {
432 ids.push(
433 cache.apply("Offline document writer", append(space, object, at, token))?,
434 );
435 at += u32::try_from(token.encode_utf16().count())?;
436 }
437 if !missing.is_empty() {
438 println!(
439 "{}",
440 json!({"event":"reviewed_append", "remote":text, "tokens":missing, "at_us":now()})
441 );
442 }
443 }
444 if generated == operations && received == ids.len() {
445 if !cache.pending()?.is_empty() {
446 return Err("Acknowledged queue did not drain".into());
447 }
448 break;
449 }
450 seed = seed
451 .wrapping_mul(6364136223846793005)
452 .wrapping_add(1442695040888963407);
453 thread::sleep(Duration::from_millis(1 + (seed >> 32) % 7));
454 }
455 Ok(())
456 })();
457 let stopped = worker.stop();
458 result?;
459 stopped?;
460 drop(cache);
461 let reopened = Replica::open(&cache_path)?;
462 if !reopened.pending()?.is_empty() {
463 return Err("Pending edits reappeared after reopen".into());
464 }
465 for id in ids {
466 let Some(EditStatus::Published { revision }) = reopened.status(id)? else {
467 return Err("Receipt did not survive reopen".into());
468 };
469 println!(
470 "{}",
471 json!({"event":"reopened_receipt", "id":id, "revision":revision.to_string()})
472 );
473 }
474 println!(
475 "{}",
476 json!({"event":"done", "operations":operations, "at_us":now()})
477 );
478 Ok(())
479}