| 1 | //! The op queue against a disposable Samba lab (`ONESTORE_SMB_LAB=127.0.0.1:PORT`, share |
| 2 | //! `agent`): publications read only the file's header, a native change is read once and |
| 3 | //! merged, and a session edits through the embedded client. |
| 4 | #![cfg(feature = "smb")] |
| 5 | |
| 6 | use notebook::{ |
| 7 | EditStatus, Remote, Replica, SmbRemote, |
| 8 | session::{Event, Section}, |
| 9 | smb::{Client, Credentials}, |
| 10 | }; |
| 11 | use onestore::{ |
| 12 | CommitError, ExGuid, Stamp, Transaction, |
| 13 | op::{Edit, Op, PageOp}, |
| 14 | }; |
| 15 | use std::{ |
| 16 | io, |
| 17 | time::{Duration, Instant, SystemTime, UNIX_EPOCH}, |
| 18 | }; |
| 19 | |
| 20 | #[path = "support/server.rs"] |
| 21 | mod server; |
| 22 | use server::snapshot; |
| 23 | |
| 24 | fn client() -> Client { |
| 25 | Client::connect( |
| 26 | &std::env::var("ONESTORE_SMB_LAB").unwrap(), |
| 27 | "agent", |
| 28 | Credentials::default(), |
| 29 | Duration::from_secs(5), |
| 30 | ) |
| 31 | .unwrap() |
| 32 | } |
| 33 | |
| 34 | fn unique(name: &str) -> String { |
| 35 | format!( |
| 36 | "{name}-{}.one", |
| 37 | SystemTime::now() |
| 38 | .duration_since(UNIX_EPOCH) |
| 39 | .unwrap() |
| 40 | .as_nanos() |
| 41 | ) |
| 42 | } |
| 43 | |
| 44 | /// The remote, counting whole-file reads. |
| 45 | struct Counted { |
| 46 | remote: SmbRemote, |
| 47 | reads: usize, |
| 48 | stamps: usize, |
| 49 | } |
| 50 | |
| 51 | impl Remote for Counted { |
| 52 | fn read(&mut self) -> io::Result<Vec<u8>> { |
| 53 | self.reads += 1; |
| 54 | self.remote.read() |
| 55 | } |
| 56 | fn stamp(&mut self) -> io::Result<Stamp> { |
| 57 | self.stamps += 1; |
| 58 | self.remote.stamp() |
| 59 | } |
| 60 | fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> { |
| 61 | self.remote.publish(transaction) |
| 62 | } |
| 63 | fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> { |
| 64 | self.remote.confirm(base) |
| 65 | } |
| 66 | } |
| 67 | |
| 68 | /// The first page's space and first body text. |
| 69 | fn target(replica: &Replica) -> (ExGuid, ExGuid) { |
| 70 | let space = replica.pages().unwrap()[0].0; |
| 71 | let text = replica |
| 72 | .page(space) |
| 73 | .unwrap() |
| 74 | .objects |
| 75 | .iter() |
| 76 | .find_map(|object| match object { |
| 77 | onestore::page::PageObject::Outline(outline) => outline |
| 78 | .paragraphs |
| 79 | .iter() |
| 80 | .find_map(|p| p.text().map(|t| t.id)), |
| 81 | _ => None, |
| 82 | }) |
| 83 | .unwrap(); |
| 84 | (space, text) |
| 85 | } |
| 86 | |
| 87 | fn typed(space: ExGuid, text: ExGuid, at: u32, with: &str) -> Edit { |
| 88 | Edit { |
| 89 | at: 133_000_000_000_000_000, |
| 90 | ops: vec![Op::Page { |
| 91 | space, |
| 92 | op: PageOp::Text { |
| 93 | text, |
| 94 | range: at..at, |
| 95 | with: with.into(), |
| 96 | }, |
| 97 | }], |
| 98 | } |
| 99 | } |
| 100 | |
| 101 | fn text_of(replica: &Replica, space: ExGuid, text: ExGuid) -> String { |
| 102 | let page = replica.page(space).unwrap(); |
| 103 | page.objects |
| 104 | .iter() |
| 105 | .find_map(|object| match object { |
| 106 | onestore::page::PageObject::Outline(outline) => outline |
| 107 | .paragraphs |
| 108 | .iter() |
| 109 | .find_map(|p| p.text().filter(|t| t.id == text)), |
| 110 | _ => None, |
| 111 | }) |
| 112 | .unwrap() |
| 113 | .text |
| 114 | .text() |
| 115 | .to_owned() |
| 116 | } |
| 117 | |
| 118 | #[test] |
| 119 | #[ignore = "requires ONESTORE_SMB_LAB pointing to disposable Samba"] |
| 120 | fn live_keystrokes_publish_reading_only_the_header_and_a_native_change_merges() { |
| 121 | let lab = client(); |
| 122 | let path = unique("queue"); |
| 123 | let source = onestore::create_section(&path, "Body", "Author").unwrap(); |
| 124 | lab.create(&path, &source).unwrap(); |
| 125 | let directory = tempfile::tempdir().unwrap(); |
| 126 | let replica = Replica::create(directory.path().join("cache.sqlite"), &source).unwrap(); |
| 127 | let (space, text) = target(&replica); |
| 128 | let mut remote = Counted { |
| 129 | remote: SmbRemote::new(client(), path.clone(), 1 << 24), |
| 130 | reads: 0, |
| 131 | stamps: 0, |
| 132 | }; |
| 133 | for n in 0..20 { |
| 134 | let at = 4 + n; |
| 135 | let id = replica |
| 136 | .apply("Author", typed(space, text, at, "x")) |
| 137 | .unwrap(); |
| 138 | assert!(matches!( |
| 139 | replica.sync_once(&mut remote).unwrap().edit, |
| 140 | Some((published, EditStatus::Published { .. })) if published == id |
| 141 | )); |
| 142 | } |
| 143 | assert_eq!(remote.reads, 0, "publications read only the header"); |
| 144 | assert_eq!(remote.stamps, 20); |
| 145 | assert_eq!(lab.read(&path, 1 << 24).unwrap(), snapshot(&replica)); |
| 146 | for _ in 0..5 { |
| 147 | assert_eq!(replica.sync_once(&mut remote).unwrap().edit, None); |
| 148 | } |
| 149 | assert_eq!(remote.reads, 0, "idle polls read only the header"); |
| 150 | |
| 151 | // Another writer types at the start of the same text; a local edit at its end merges. |
| 152 | let native = { |
| 153 | let arena = onestore::Arena::default(); |
| 154 | let mut section = |
| 155 | onestore::Section::open(&arena, lab.read(&path, 1 << 24).unwrap()).unwrap(); |
| 156 | section |
| 157 | .apply("Native", &typed(space, text, 0, "Native ")) |
| 158 | .unwrap(); |
| 159 | section.seal().unwrap().unwrap() |
| 160 | }; |
| 161 | lab.commit_transaction(&path, &native).unwrap(); |
| 162 | let end = text_of(&replica, space, text).encode_utf16().count() as u32; |
| 163 | let id = replica |
| 164 | .apply("Author", typed(space, text, end, " local")) |
| 165 | .unwrap(); |
| 166 | let synced = replica.sync_once(&mut remote).unwrap(); |
| 167 | assert_eq!(synced.changed, [space]); |
| 168 | assert!( |
| 169 | matches!(synced.edit, Some((published, EditStatus::Published { .. })) if published == id) |
| 170 | ); |
| 171 | assert_eq!(remote.reads, 1, "a changed remote is read once"); |
| 172 | let merged = text_of(&replica, space, text); |
| 173 | assert!(merged.starts_with("Native Body"), "{merged}"); |
| 174 | assert!(merged.ends_with(" local"), "{merged}"); |
| 175 | assert_eq!(lab.read(&path, 1 << 24).unwrap(), snapshot(&replica)); |
| 176 | lab.delete(&path).unwrap(); |
| 177 | } |
| 178 | |
| 179 | #[test] |
| 180 | #[ignore = "requires ONESTORE_SMB_LAB pointing to disposable Samba"] |
| 181 | fn live_session_applies_edits_through_the_embedded_client_and_reopens() { |
| 182 | let lab = client(); |
| 183 | let path = unique("session"); |
| 184 | let source = onestore::create_section(&path, "Session", "Author").unwrap(); |
| 185 | lab.create(&path, &source).unwrap(); |
| 186 | let directory = tempfile::tempdir().unwrap(); |
| 187 | let cache = directory.path().join("cache.sqlite"); |
| 188 | let section = Section::resume_smb( |
| 189 | path.clone(), |
| 190 | Replica::create(&cache, &source).unwrap(), |
| 191 | 1 << 24, |
| 192 | || { |
| 193 | Client::connect( |
| 194 | &std::env::var("ONESTORE_SMB_LAB").unwrap(), |
| 195 | "agent", |
| 196 | Credentials::default(), |
| 197 | Duration::from_secs(5), |
| 198 | ) |
| 199 | }, |
| 200 | || {}, |
| 201 | ) |
| 202 | .unwrap(); |
| 203 | let (space, text) = target(section.replica()); |
| 204 | for n in 0..10 { |
| 205 | section |
| 206 | .apply("Author", typed(space, text, n, &n.to_string())) |
| 207 | .unwrap(); |
| 208 | } |
| 209 | // Applying does not wait: the edits are done when the file holds them. |
| 210 | let remote_text = |bytes: Vec<u8>| { |
| 211 | let arena = onestore::Arena::default(); |
| 212 | let section = onestore::Section::open(&arena, bytes).unwrap(); |
| 213 | let page = section.page(space).unwrap(); |
| 214 | page.objects |
| 215 | .iter() |
| 216 | .find_map(|object| match object { |
| 217 | onestore::page::PageObject::Outline(outline) => { |
| 218 | outline.paragraphs.iter().find_map(|p| { |
| 219 | p.text() |
| 220 | .filter(|t| t.id == text) |
| 221 | .map(|t| t.text.text().to_owned()) |
| 222 | }) |
| 223 | } |
| 224 | _ => None, |
| 225 | }) |
| 226 | .unwrap() |
| 227 | }; |
| 228 | let deadline = Instant::now() + Duration::from_secs(30); |
| 229 | while remote_text(lab.read(&path, 1 << 24).unwrap()) != "0123456789Session" |
| 230 | || !section.pending().unwrap().is_empty() |
| 231 | { |
| 232 | assert!(Instant::now() < deadline, "the edits were not published"); |
| 233 | for event in section.events() { |
| 234 | assert!( |
| 235 | !matches!(event, Event::Rejected { .. } | Event::Failed(_)), |
| 236 | "{event:?}" |
| 237 | ); |
| 238 | } |
| 239 | std::thread::sleep(Duration::from_millis(50)); |
| 240 | } |
| 241 | let published = lab.read(&path, 1 << 24).unwrap(); |
| 242 | section.close().unwrap(); |
| 243 | let replica = Replica::open(&cache).unwrap(); |
| 244 | assert_eq!(text_of(&replica, space, text), "0123456789Session"); |
| 245 | assert_eq!(snapshot(&replica), published); |
| 246 | lab.delete(&path).unwrap(); |
| 247 | } |