| ... | @@ -57,6 +57,9 @@ fn retry(message: &str) -> Error { | ... | @@ -57,6 +57,9 @@ fn retry(message: &str) -> Error { |
| 57 | impl Served { | 57 | impl Served { |
| 58 | pub(super) fn batch(self: &Arc<Self>, peer: &[u8; 16], request: Request) -> Result<Reply> { | 58 | pub(super) fn batch(self: &Arc<Self>, peer: &[u8; 16], request: Request) -> Result<Reply> { |
| 59 | let path = request.path.clone(); | 59 | let path = request.path.clone(); |
| | 60 | if !self.writers.lock().unwrap().contains_key(&path) { |
| | 61 | self.image(&path)?; |
| | 62 | } |
| 60 | let mut writers = self.writers.lock().unwrap(); | 63 | let mut writers = self.writers.lock().unwrap(); |
| 61 | let writer = writers.entry(path.clone()).or_insert_with(|| { | 64 | let writer = writers.entry(path.clone()).or_insert_with(|| { |
| 62 | let (send, receive) = mpsc::sync_channel::<Waiting>(64); | 65 | let (send, receive) = mpsc::sync_channel::<Waiting>(64); |
| ... | @@ -302,6 +305,121 @@ mod tests { | ... | @@ -302,6 +305,121 @@ mod tests { |
| 302 | page::{Attachment, PageObject}, | 305 | page::{Attachment, PageObject}, |
| 303 | }; | 306 | }; |
| 304 | | 307 | |
| | 308 | #[cfg(not(target_arch = "wasm32"))] |
| | 309 | #[test] |
| | 310 | fn missing_sections_do_not_retain_writers_and_can_be_created_later() { |
| | 311 | let directory = tempfile::tempdir().unwrap(); |
| | 312 | let folder = directory.path().join("notebook"); |
| | 313 | std::fs::create_dir(&folder).unwrap(); |
| | 314 | let served = Arc::new(Served { |
| | 315 | storage: crate::session::Notebook::open(&folder, directory.path().join("cache")) |
| | 316 | .unwrap() |
| | 317 | .into_storage(), |
| | 318 | images: Mutex::default(), |
| | 319 | snapshots: Mutex::default(), |
| | 320 | puts: Mutex::default(), |
| | 321 | guests: Mutex::default(), |
| | 322 | writers: Mutex::default(), |
| | 323 | host: Mutex::default(), |
| | 324 | room: Mutex::default(), |
| | 325 | }); |
| | 326 | let peer = [1; 16]; |
| | 327 | let (line, _outbound) = mpsc::channel(); |
| | 328 | let (queue, _requests) = mpsc::sync_channel(1); |
| | 329 | served.guests.lock().unwrap().insert( |
| | 330 | peer, |
| | 331 | Admitted { |
| | 332 | line: crate::live::Line(line), |
| | 333 | queue, |
| | 334 | starts: BURST, |
| | 335 | counted: Instant::now(), |
| | 336 | held: Vec::new(), |
| | 337 | }, |
| | 338 | ); |
| | 339 | for at in 0..128 { |
| | 340 | let reply = served.handle( |
| | 341 | &peer, |
| | 342 | kind::EDITS, |
| | 343 | &minicbor::to_vec(Request { |
| | 344 | path: format!("Missing{at}.one"), |
| | 345 | ..Request::default() |
| | 346 | }) |
| | 347 | .unwrap(), |
| | 348 | ); |
| | 349 | assert!(reply.failure.is_some()); |
| | 350 | } |
| | 351 | assert!(served.writers.lock().unwrap().is_empty()); |
| | 352 | |
| | 353 | let path = "Missing0.one"; |
| | 354 | let image = onestore::create_section(path, "Original text", "Fixture").unwrap(); |
| | 355 | std::fs::write(folder.join(path), &image).unwrap(); |
| | 356 | let store = Store::parse(&image).unwrap(); |
| | 357 | let index = RevisionIndex::parse(&store).unwrap(); |
| | 358 | let document = onestore::document::Document::parse(&index).unwrap(); |
| | 359 | let (space, text) = document |
| | 360 | .spaces |
| | 361 | .iter() |
| | 362 | .find_map(|(space, object)| { |
| | 363 | object.revisions[&object.contexts[&ExGuid::default()]] |
| | 364 | .nodes |
| | 365 | .iter() |
| | 366 | .find_map(|(id, node)| { |
| | 367 | matches!(node.kind, onestore::document::Kind::RichText { .. }) |
| | 368 | .then_some((*space, *id)) |
| | 369 | }) |
| | 370 | }) |
| | 371 | .unwrap(); |
| | 372 | let edit = Edit { |
| | 373 | at: 134_000_000_000_000_000, |
| | 374 | ops: vec![Op::Page { |
| | 375 | space, |
| | 376 | op: PageOp::Text { |
| | 377 | text, |
| | 378 | range: 13..13, |
| | 379 | with: "X".into(), |
| | 380 | }, |
| | 381 | }], |
| | 382 | }; |
| | 383 | let arena = Arena::default(); |
| | 384 | let mut section = Section::open(&arena, image).unwrap(); |
| | 385 | section.apply("Alice", &edit).unwrap(); |
| | 386 | let transaction = section.seal().unwrap().unwrap(); |
| | 387 | let request = Request { |
| | 388 | path: path.into(), |
| | 389 | stamp: Some(transaction.base().into()), |
| | 390 | bytes: Some( |
| | 391 | encode( |
| | 392 | &Edits { |
| | 393 | edits: vec![("Alice".into(), edit)], |
| | 394 | revisions: section.newest().filter(|(at, _)| *at == space).collect(), |
| | 395 | }, |
| | 396 | LIMIT, |
| | 397 | ) |
| | 398 | .unwrap(), |
| | 399 | ), |
| | 400 | ..Request::default() |
| | 401 | }; |
| | 402 | let first = served.batch(&peer, request.clone()).unwrap(); |
| | 403 | assert!(first.failure.is_none(), "{:?}", first.failure); |
| | 404 | let committed = std::fs::read(folder.join(path)).unwrap(); |
| | 405 | let stamp = Stamp::of(&committed).unwrap(); |
| | 406 | assert_eq!( |
| | 407 | Stamp::try_from(first.stamp.as_ref().unwrap()).unwrap(), |
| | 408 | stamp |
| | 409 | ); |
| | 410 | let read_arena = Arena::default(); |
| | 411 | let reread = Section::open(&read_arena, committed.clone()).unwrap(); |
| | 412 | assert_eq!(reread.page(space).unwrap(), section.page(space).unwrap()); |
| | 413 | let repeated = served.batch(&peer, request).unwrap(); |
| | 414 | assert!(repeated.failure.is_none(), "{:?}", repeated.failure); |
| | 415 | let replayed = std::fs::read(folder.join(path)).unwrap(); |
| | 416 | assert_eq!(&replayed[1024..], &committed[1024..]); |
| | 417 | let replay_arena = Arena::default(); |
| | 418 | let replay = Section::open(&replay_arena, replayed).unwrap(); |
| | 419 | assert_eq!(replay.page(space).unwrap(), section.page(space).unwrap()); |
| | 420 | assert_eq!(served.writers.lock().unwrap().len(), 1); |
| | 421 | } |
| | 422 | |
| 305 | #[test] | 423 | #[test] |
| 306 | fn attachment_encoding_stops_at_the_upload_budget() { | 424 | fn attachment_encoding_stops_at_the_upload_budget() { |
| 307 | let id = ExGuid { | 425 | let id = ExGuid { |