| 1 | #[path = "../../onestore/examples/support/concurrent.rs"] |
| 2 | mod concurrent; |
| 3 | #[path = "support/edit.rs"] |
| 4 | mod edit; |
| 5 | |
| 6 | use notebook::smb::{Client, Credentials}; |
| 7 | use onestore::{ |
| 8 | CommitError, CommitState, ExGuid, RevisionIndex, Stamp, Store, |
| 9 | document::{Document, Kind}, |
| 10 | }; |
| 11 | use serde_json::json; |
| 12 | use std::{ |
| 13 | cell::RefCell, |
| 14 | env, io, thread, |
| 15 | time::{Duration, Instant, SystemTime, UNIX_EPOCH}, |
| 16 | }; |
| 17 | |
| 18 | fn 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. |
| 36 | fn 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 | |
| 52 | fn 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] |
| 251 | fn 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 | } |