1use super::*;
2use crate::{PendingEdit, merge};
3use onestore::{Arena, ExGuid, Section, op::Edit};
4use std::collections::BTreeSet;
5
6#[derive(serde::Serialize, serde::Deserialize)]
7#[serde(deny_unknown_fields)]
8pub(super) struct Edits {
9 pub edits: Vec<(String, Edit)>,
10 pub revisions: BTreeMap<ExGuid, ExGuid>,
11}
12
13pub(super) fn encode(edits: &Edits, limit: usize) -> io::Result<Vec<u8>> {
14 struct Buffer {
15 bytes: Vec<u8>,
16 limit: usize,
17 }
18 impl io::Write for Buffer {
19 fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
20 if bytes.len() > self.limit.saturating_sub(self.bytes.len()) {
21 return Err(io::ErrorKind::FileTooLarge.into());
22 }
23 self.bytes.extend_from_slice(bytes);
24 Ok(bytes.len())
25 }
26
27 fn flush(&mut self) -> io::Result<()> {
28 Ok(())
29 }
30 }
31 let mut buffer = Buffer {
32 bytes: Vec::new(),
33 limit,
34 };
35 serde_json::to_writer(&mut buffer, edits).map_err(|error| {
36 io::Error::new(
37 error.io_error_kind().unwrap_or(io::ErrorKind::InvalidData),
38 error,
39 )
40 })?;
41 Ok(buffer.bytes)
42}
43
44pub(super) struct Waiting {
45 peer: [u8; 16],
46 request: Request,
47 reply: mpsc::Sender<Reply>,
48}
49
50fn retry(message: &str) -> Error {
51 Error::Remote(CommitError {
52 state: CommitState::NotCommitted,
53 error: io::Error::new(io::ErrorKind::ResourceBusy, message.to_owned()),
54 })
55}
56
57impl Served {
58 pub(super) fn batch(self: &Arc<Self>, peer: &[u8; 16], request: Request) -> Result<Reply> {
59 let path = request.path.clone();
60 let mut writers = self.writers.lock().unwrap();
61 let writer = writers.entry(path.clone()).or_insert_with(|| {
62 let (send, receive) = mpsc::sync_channel::<Waiting>(64);
63 let served = Arc::downgrade(self);
64 thread::spawn(move || {
65 while let Ok(first) = receive.recv() {
66 let mut waiting = vec![first];
67 if let Ok(next) = receive.recv_timeout(Duration::from_millis(5)) {
68 waiting.push(next);
69 }
70 waiting.extend(receive.try_iter().take(62));
71 let Some(served) = served.upgrade() else {
72 return;
73 };
74 served.write_batches(&path, waiting);
75 }
76 });
77 send
78 });
79 let (reply, receive) = mpsc::channel();
80 writer
81 .try_send(Waiting {
82 peer: *peer,
83 request,
84 reply,
85 })
86 .map_err(|_| retry("The section's writer is busy"))?;
87 drop(writers);
88 receive.recv_timeout(TIMEOUT).map_err(|_| {
89 Error::Remote(CommitError {
90 state: CommitState::Unknown,
91 error: io::Error::new(
92 io::ErrorKind::TimedOut,
93 "The section's writer did not answer",
94 ),
95 })
96 })
97 }
98
99 fn write_batches(&self, path: &str, waiting: Vec<Waiting>) {
100 let before = match self.image(path) {
101 Ok(image) => image,
102 Err(error) => {
103 for batch in waiting {
104 let _ = batch.reply.send(failed(0, retry(&error.to_string())));
105 }
106 return;
107 }
108 };
109 let mut image = Arc::clone(&before);
110 let mut combined: Option<Transaction> = None;
111 let mut accepted = Vec::new();
112 for batch in waiting {
113 match self.prepare_batch(path, &image, &batch) {
114 Ok(Some(transaction)) => {
115 let mut next = (*image).clone();
116 let extended =
117 transaction
118 .apply(&mut next)
119 .and_then(|()| match &mut combined {
120 Some(combined) => combined.extend(transaction),
121 None => {
122 combined = Some(transaction);
123 Ok(())
124 }
125 });
126 match extended {
127 Ok(()) => {
128 image = Arc::new(next);
129 accepted.push(batch);
130 }
131 Err(error) => {
132 let _ = batch.reply.send(failed(0, error.into()));
133 }
134 }
135 }
136 Ok(None) => accepted.push(batch),
137 Err(error) => {
138 let _ = batch.reply.send(failed(0, error));
139 }
140 }
141 }
142 if accepted.is_empty() {
143 return;
144 }
145 let result = match &combined {
146 Some(transaction) => self.storage.commit(path, transaction),
147 None => self
148 .storage
149 .confirm(path, &Stamp::of(&image).expect("section stamp"))
150 .map_err(Error::from),
151 };
152 if result.is_ok() {
153 self.tell_delta(path, &before, &image, None);
154 self.keep(path, Arc::clone(&image));
155 self.changed(&[path.to_owned()]);
156 }
157 for batch in accepted {
158 let reply = match &result {
159 Ok(()) => Reply {
160 stamp: Some((&Stamp::of(&image).expect("section stamp")).into()),
161 ..Reply::default()
162 },
163 Err(Error::Remote(error)) => failed(
164 0,
165 Error::Remote(CommitError {
166 state: error.state,
167 error: io::Error::new(error.error.kind(), error.error.to_string()),
168 }),
169 ),
170 Err(error) => failed(0, retry(&error.to_string())),
171 };
172 let _ = batch.reply.send(reply);
173 }
174 }
175
176 fn prepare_batch(
177 &self,
178 path: &str,
179 image: &[u8],
180 batch: &Waiting,
181 ) -> Result<Option<Transaction>> {
182 if !self.guests.lock().unwrap().contains_key(&batch.peer) {
183 return Err(refused(
184 io::ErrorKind::PermissionDenied,
185 "The device is no longer connected",
186 ));
187 }
188 let stamp: Stamp = batch
189 .request
190 .stamp
191 .as_ref()
192 .ok_or_else(|| retry("No batch base"))?
193 .try_into()?;
194 let edits: Edits = serde_json::from_slice(&self.carried(&batch.peer, &batch.request)?)
195 .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?;
196 if edits.revisions.is_empty()
197 || edits
198 .revisions
199 .values()
200 .any(|revision| revision.guid == [0; 16])
201 {
202 return Err(retry("No batch revisions"));
203 }
204 let store = Store::parse(image)?;
205 let index = RevisionIndex::parse(&store)?;
206 if edits.revisions.iter().all(|(space, revision)| {
207 index
208 .spaces
209 .get(space)
210 .is_some_and(|space| space.revisions.contains_key(revision))
211 }) {
212 return Ok(None);
213 }
214 let names = edits
215 .revisions
216 .iter()
217 .filter(|(space, revision)| {
218 !index
219 .spaces
220 .get(space)
221 .is_some_and(|space| space.revisions.contains_key(revision))
222 })
223 .map(|(space, revision)| (*space, *revision))
224 .collect();
225 let arena = Arena::default();
226 let mut next = Section::open(&arena, image.to_vec())?;
227 let queued: Vec<PendingEdit> = edits
228 .edits
229 .into_iter()
230 .enumerate()
231 .map(|(id, (author, edit))| PendingEdit {
232 id: id as u64,
233 author,
234 edit,
235 })
236 .collect();
237 if Stamp::of(image)? == stamp {
238 for edit in &queued {
239 next.apply(&edit.author, &edit.edit)?;
240 }
241 } else {
242 let base = self
243 .images
244 .lock()
245 .unwrap()
246 .iter()
247 .find_map(|(held, at, image)| {
248 (held == path && *at == stamp).then(|| Arc::clone(image))
249 })
250 .ok_or_else(|| retry("The batch's base is no longer held"))?;
251 let old_arena = Arena::default();
252 let mut old = Section::open(&old_arena, (*base).clone())?;
253 if old.root() != next.root() {
254 return Err(retry("The batch belongs to another section"));
255 }
256 let merged = merge::rebase(&mut old, &mut next, &queued, &BTreeSet::new())?;
257 if !merged.conflicts.is_empty() {
258 let local_arena = Arena::default();
259 let mut local = Section::open(&local_arena, (*base).clone())?;
260 for edit in &queued {
261 local.apply(&edit.author, &edit.edit)?;
262 }
263 for (space, (author, objects)) in &merged.conflicts {
264 let page = local.page(*space)?;
265 if next.page(*space).is_ok_and(|remote| remote == page) {
266 return Err(retry("The batch's edits already landed"));
267 }
268 merge::conflict_page(
269 &mut next,
270 &mut local,
271 &merged.moved,
272 *space,
273 &page,
274 author,
275 objects,
276 crate::now(),
277 None,
278 )?;
279 }
280 }
281 }
282 let transaction = next.seal_as(&names)?;
283 if edits.revisions.iter().any(|(space, revision)| {
284 !next
285 .newest()
286 .any(|(held, at)| held == *space && at == *revision)
287 }) {
288 return Err(refused(
289 io::ErrorKind::Unsupported,
290 "The batch needs a guest's commit",
291 ));
292 }
293 Ok(transaction)
294 }
295}
296
297#[cfg(test)]
298mod tests {
299 use super::*;
300 use onestore::{
301 op::{Op, PageOp},
302 page::{Attachment, PageObject},
303 };
304
305 #[test]
306 fn attachment_encoding_stops_at_the_upload_budget() {
307 let id = ExGuid {
308 guid: [1; 16],
309 n: 1,
310 };
311 let bytes: Arc<[u8]> = (0..4096)
312 .map(|at| (at % 251) as u8)
313 .collect::<Vec<_>>()
314 .into();
315 let edits = Edits {
316 edits: vec![(
317 "Ada".into(),
318 Edit {
319 at: 0,
320 ops: vec![Op::Page {
321 space: id,
322 op: PageOp::Add {
323 object: PageObject::Attachment(Attachment {
324 id,
325 filename: "payload.bin".into(),
326 source_path: None,
327 size: None,
328 layout: Default::default(),
329 bytes: Some(Arc::clone(&bytes)),
330 preview: None,
331 recording: None,
332 tags: Vec::new(),
333 }),
334 before: None,
335 },
336 }],
337 },
338 )],
339 revisions: BTreeMap::new(),
340 };
341 let encoded = serde_json::to_vec(&edits).unwrap();
342 assert_eq!(encode(&edits, encoded.len()).unwrap(), encoded);
343 let decoded: Edits = serde_json::from_slice(&encoded).unwrap();
344 assert_eq!(decoded.edits, edits.edits);
345 assert_eq!(
346 encode(&edits, encoded.len() - 1).unwrap_err().kind(),
347 io::ErrorKind::FileTooLarge
348 );
349 assert_eq!(
350 encode(&edits, 128).unwrap_err().kind(),
351 io::ErrorKind::FileTooLarge
352 );
353 }
354}