1#[path = "../../onestore/examples/support/concurrent.rs"]
2mod concurrent;
3#[path = "support/edit.rs"]
4mod edit;
5
6use notebook::smb::{Client, Credentials};
7use onestore::{
8 CommitError, CommitState, ExGuid, RevisionIndex, Stamp, Store,
9 document::{Document, Kind},
10};
11use serde_json::json;
12use std::{
13 cell::RefCell,
14 env, io, thread,
15 time::{Duration, Instant, SystemTime, UNIX_EPOCH},
16};
17
18fn paragraph(
19 bytes: &[u8],
20 space: ExGuid,
21 object: ExGuid,
22) -> Result<String, Box<dyn std::error::Error>> {
23 let store = Store::parse(bytes)?;
24 let index = RevisionIndex::parse(&store)?;
25 index.validate_current()?;
26 let document = Document::parse(&index)?;
27 let space = &document.spaces[&space];
28 let revision = &space.revisions[&space.contexts[&ExGuid::default()]];
29 let Kind::RichText { text, .. } = &revision.nodes[&object].kind else {
30 return Err("The append target is no longer text.".into());
31 };
32 Ok(text.clone())
33}
34
35// Only the owned append workload guarantees unique tokens that no writer removes.
36fn retained(before: &str, token: &str, current: &str, state: CommitState) -> io::Result<bool> {
37 let count = current.matches(token).count();
38 if count == 0 && state != CommitState::Committed && current.starts_with(before) {
39 return Ok(false);
40 }
41 if count == 1
42 && state != CommitState::NotCommitted
43 && current.starts_with(&format!("{before}{token}"))
44 {
45 return Ok(true);
46 }
47 Err(io::Error::other(
48 "The append history contradicts the commit outcome.",
49 ))
50}
51
52fn main() -> Result<(), Box<dyn std::error::Error>> {
53 let args: Vec<_> = env::args().skip(1).collect();
54 if args
55 .first()
56 .is_none_or(|mode| !["read", "write"].contains(&mode.as_str()))
57 {
58 return Err("Reconnect testing requires the append-only workload.".into());
59 }
60 let address = env::var("ONESTORE_SMB_LAB")?;
61 let share = env::var("ONESTORE_SMB_SHARE")?;
62 let client = RefCell::new(Some(Client::connect(
63 &address,
64 &share,
65 Credentials::default(),
66 Duration::from_secs(5),
67 )?));
68 let read = |path: &str| {
69 let mut session = client.borrow_mut();
70 if session.is_none() {
71 match Client::connect(
72 &address,
73 &share,
74 Credentials::default(),
75 Duration::from_secs(5),
76 ) {
77 Ok(fresh) => {
78 *session = Some(fresh);
79 println!(
80 "{}",
81 json!({"event":"transport_connected", "at_us":SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_micros()})
82 );
83 }
84 Err(error) => {
85 println!(
86 "{}",
87 json!({"event":"transport_connect_error","error":error.to_string()})
88 );
89 return Err(io::ErrorKind::WouldBlock.into());
90 }
91 }
92 }
93 let started = SystemTime::now()
94 .duration_since(UNIX_EPOCH)
95 .unwrap()
96 .as_micros();
97 match session.as_ref().unwrap().read(path, 256 * 1024 * 1024) {
98 Ok(bytes) => {
99 if args.get(2).is_some_and(|actor| actor == "r0")
100 && let Some(folder) = env::var_os("ONESTORE_OFFLINE_FORMAT_REPLY_DIR")
101 {
102 let folder = std::path::PathBuf::from(folder);
103 let captured = folder.join("offline-retired.one");
104 if !captured.exists()
105 && let Ok(marker) = std::fs::read(folder.join("offline-paused-w0.isolate"))
106 && let Ok(watched) = serde_json::from_slice::<serde_json::Value>(&marker)
107 && watched["after_us"]
108 .as_u64()
109 .is_some_and(|after| started > u128::from(after))
110 {
111 let sid: ExGuid = watched["space"]
112 .as_str()
113 .ok_or_else(|| io::Error::other("Missing watched space"))?
114 .parse()
115 .map_err(io::Error::other)?;
116 let rid: ExGuid = watched["revision"]
117 .as_str()
118 .ok_or_else(|| io::Error::other("Missing watched revision"))?
119 .parse()
120 .map_err(io::Error::other)?;
121 let store = Store::parse(&bytes).map_err(io::Error::other)?;
122 let index = RevisionIndex::parse(&store).map_err(io::Error::other)?;
123 if index
124 .spaces
125 .get(&sid)
126 .is_some_and(|space| !space.revisions.contains_key(&rid))
127 {
128 index.validate_current().map_err(io::Error::other)?;
129 std::fs::write(captured.with_extension("tmp"), &bytes)?;
130 std::fs::rename(captured.with_extension("tmp"), captured)?;
131 println!(
132 "{}",
133 json!({"event":"revision_retired", "space":sid.to_string(), "revision":rid.to_string(), "started_us":started, "finished_us":SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_micros()})
134 );
135 }
136 }
137 }
138 Ok(bytes)
139 }
140 Err(error)
141 if matches!(
142 error.kind(),
143 io::ErrorKind::Other | io::ErrorKind::TimedOut | io::ErrorKind::NotConnected
144 ) =>
145 {
146 println!(
147 "{}",
148 json!({"event":"transport_read_error","error":error.to_string()})
149 );
150 *session = None;
151 Err(io::ErrorKind::WouldBlock.into())
152 }
153 result => result,
154 }
155 };
156 concurrent::run(
157 &args,
158 read,
159 |path, source, space, object, range, replacement| {
160 let outcome =
161 edit::replaced(source, space, object, range, replacement).and_then(|transaction| {
162 client
163 .borrow()
164 .as_ref()
165 .unwrap()
166 .commit_transaction(path, &transaction)
167 });
168 let Err(error) = outcome else {
169 return Ok(());
170 };
171 if error.state == CommitState::NotCommitted
172 && matches!(
173 error.error.kind(),
174 io::ErrorKind::WouldBlock
175 | io::ErrorKind::ResourceBusy
176 | io::ErrorKind::PermissionDenied
177 | io::ErrorKind::NotFound
178 )
179 {
180 return Err(error);
181 }
182 println!(
183 "{}",
184 json!({"event":"transport_commit_error","state":format!("{:?}",error.state),"token":replacement,"error":error.error.to_string()})
185 );
186 *client.borrow_mut() = None;
187 let deadline = Instant::now() + Duration::from_secs(60);
188 loop {
189 if Instant::now() >= deadline {
190 return Err(error);
191 }
192 let current = match read(path) {
193 Ok(bytes) => bytes,
194 Err(retry) if retry.kind() == io::ErrorKind::WouldBlock => {
195 thread::sleep(Duration::from_millis(100));
196 continue;
197 }
198 Err(_) => return Err(error),
199 };
200 let published = paragraph(source, space, object)
201 .and_then(|before| {
202 Ok(retained(
203 &before,
204 replacement,
205 &paragraph(&current, space, object)?,
206 error.state,
207 )?)
208 })
209 .map_err(|failure| CommitError {
210 state: error.state,
211 error: io::Error::other(failure.to_string()),
212 })?;
213 if published {
214 let confirmation = Stamp::of(&current)
215 .map_err(|failure| CommitError {
216 state: CommitState::NotCommitted,
217 error: io::Error::other(failure.to_string()),
218 })
219 .and_then(|base| client.borrow().as_ref().unwrap().confirm(path, &base));
220 if let Err(failure) = confirmation
221 && failure.state != CommitState::Committed
222 {
223 println!(
224 "{}",
225 json!({"event":"transport_confirmation_error","state":format!("{:?}", failure.state),"error":failure.error.to_string()})
226 );
227 *client.borrow_mut() = None;
228 thread::sleep(Duration::from_millis(100));
229 continue;
230 }
231 println!(
232 "{}",
233 json!({"event":"transport_reconciled","token":replacement,"published":true,"flush_confirmed":true})
234 );
235 return Ok(());
236 }
237 println!(
238 "{}",
239 json!({"event":"transport_reconciled","token":replacement,"published":false})
240 );
241 return Err(CommitError {
242 state: CommitState::NotCommitted,
243 error: io::ErrorKind::ResourceBusy.into(),
244 });
245 }
246 },
247 )
248}
249
250#[test]
251fn uncertain_append_requires_one_retained_token_and_its_predecessor() {
252 for state in [CommitState::Unknown, CommitState::Committed] {
253 assert!(retained("before", " [w0:0]", "before [w0:0] [w1:0]", state).unwrap());
254 }
255 for state in [CommitState::Unknown, CommitState::NotCommitted] {
256 assert!(!retained("before", " [w0:0]", "before [w1:0]", state).unwrap());
257 }
258 for (current, state) in [
259 ("before [w0:0] [w0:0]", CommitState::Unknown),
260 ("changed [w0:0]", CommitState::Unknown),
261 ("befor", CommitState::Unknown),
262 ("before", CommitState::Committed),
263 ("before [w0:0]", CommitState::NotCommitted),
264 ] {
265 assert!(retained("before", " [w0:0]", current, state).is_err());
266 }
267}