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
6use notebook::{
7 EditStatus, Remote, Replica, SmbRemote,
8 session::{Event, Section},
9 smb::{Client, Credentials},
10};
11use onestore::{
12 CommitError, ExGuid, Stamp, Transaction,
13 op::{Edit, Op, PageOp},
14};
15use std::{
16 io,
17 time::{Duration, Instant, SystemTime, UNIX_EPOCH},
18};
19
20#[path = "support/server.rs"]
21mod server;
22use server::snapshot;
23
24fn 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
34fn 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.
45struct Counted {
46 remote: SmbRemote,
47 reads: usize,
48 stamps: usize,
49}
50
51impl 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.
69fn 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
87fn 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
101fn 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"]
120fn 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"]
181fn 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}