1use std::{
2 collections::{HashMap, HashSet},
3 num::NonZeroUsize,
4 sync::Arc,
5 thread,
6 time::{Duration, Instant},
7};
8
9use futures::{StreamExt, stream};
10use tokio::sync::{Mutex, broadcast::error::TryRecvError};
11use tokio::time;
12use usvg::Transform;
13use uuid::Uuid;
14
15use crate::{
16 LbErrKind, LbResult, LocalLb,
17 io::network::ApiError,
18 model::{
19 ValidationFailure,
20 access_info::UserAccessMode,
21 account::Account,
22 api::{
23 ChangeDocRequestV2, GetDocRequest, GetFileIdsRequest, GetUpdatesRequestV2,
24 GetUsernameError, GetUsernameRequest, UpsertDebugInfoRequest, UpsertRequestV2,
25 },
26 chat,
27 crypto::{DecryptedDocument, EncryptedDocument},
28 errors::{LbErr, Unexpected},
29 file::ShareMode,
30 file_like::FileLike,
31 file_metadata::{DocumentHmac, FileDiff, FileType, Owner},
32 filename::{DocumentType, NameComponents},
33 lazy::LazyTree,
34 signed_meta::SignedMeta,
35 staged::StagedTreeLikeMut,
36 svg::{self, buffer::u_transform_to_bezier, element::Element},
37 symkey, text,
38 tree_like::TreeLike,
39 validate,
40 },
41 service::events::{Actor, Event, SyncIncrement},
42};
43
44pub type Syncer = Arc<Mutex<SyncState>>;
45
46#[derive(Default)]
47pub struct SyncState {
48 last_synced: u64,
50
51 updates_as_of: u64,
53
54 remote_changes: Vec<SignedMeta>,
56
57 new_root: Option<Uuid>,
59
60 pulled_docs: Vec<Uuid>,
62}
63
64impl LocalLb {
71 #[instrument(level = "debug", skip(self), err(Debug))]
72 pub async fn sync(&self) -> LbResult<()> {
73 let mut sync_state = self.syncer.lock().await;
74 self.events.sync_update(SyncIncrement::SyncStarted);
75
76 let pipeline: LbResult<()> = async {
77 self.pull_updates(&mut sync_state).await?;
78 self.push_local_changes().await?;
79 Ok(())
80 }
81 .await;
82
83 self.events.sync_update(SyncIncrement::SyncFinished(
84 pipeline.as_ref().err().map(|err| err.kind.clone()),
85 ));
86
87 self.cleanup().await?;
88
89 pipeline?;
90
91 let account = self.get_account()?.clone();
92
93 #[cfg(not(target_family = "wasm"))]
94 if account.is_beta() {
95 self.send_debug_info(account).await;
96 }
97
98 Ok(())
99 }
100
101 pub(crate) async fn pull_updates(&self, sync_state: &mut SyncState) -> LbResult<()> {
102 self.inital_sync_state(sync_state).await?;
103 self.process_deletions().await?;
104 self.fetch_meta(sync_state).await?;
105 self.fetch_required_docs(sync_state).await?;
106 self.merge(sync_state).await?;
108 self.commit_last_synced(sync_state).await?;
109 self.send_pull_events(sync_state).await?;
110
111 if !self.config.background_work {
112 self.populate_pk_cache().await?;
113 }
114
115 Ok(())
116 }
117
118 pub(crate) async fn push_local_changes(&self) -> LbResult<()> {
119 self.push_meta().await?;
120 self.push_docs().await?;
121
122 Ok(())
123 }
124
125 async fn inital_sync_state(&self, state: &mut SyncState) -> LbResult<()> {
126 let tx = self.ro_tx().await;
127 let db = tx.db();
128
129 *state = Default::default();
130 state.last_synced = db.last_synced.get().copied().unwrap_or_default() as u64;
131
132 Ok(())
133 }
134
135 pub(crate) async fn process_deletions(&self) -> LbResult<()> {
136 let server_ids = self
137 .client
138 .request(self.get_account()?, GetFileIdsRequest {})
139 .await?
140 .ids;
141
142 let mut tx = self.begin_tx().await;
143 let db = tx.db();
144
145 let mut local = db.base_metadata.stage(&db.local_metadata).to_lazy();
146 let base_ids = local.tree.base.ids();
147
148 let mut prunable_ids = base_ids;
149 prunable_ids.retain(|id| !server_ids.contains(id));
150 for id in prunable_ids.clone() {
151 prunable_ids.extend(local.descendants(&id)?.into_iter());
152 }
153 for id in &prunable_ids {
154 if let Some(base_file) = local.tree.base.maybe_find(id) {
155 self.docs
156 .delete(*id, base_file.document_hmac().copied())
157 .await?;
158 }
159 if let Some(local_file) = local.maybe_find(id) {
160 self.docs
161 .delete(*id, local_file.document_hmac().copied())
162 .await?;
163 }
164 }
165
166 let mut base_staged = (&mut db.base_metadata).to_lazy().stage(None);
167 base_staged.tree.removed = prunable_ids.iter().copied().collect();
168 base_staged.promote()?;
169
170 let mut local_staged = (&mut db.local_metadata).to_lazy().stage(None);
171 local_staged.tree.removed = prunable_ids.iter().copied().collect();
172 local_staged.promote()?;
173
174 if !prunable_ids.is_empty() {
175 self.events.meta_changed(Actor::Sync);
176 }
177
178 Ok(())
179 }
180
181 async fn fetch_meta(&self, state: &mut SyncState) -> LbResult<()> {
182 let updates = self
183 .client
184 .request(
185 self.get_account()?,
186 GetUpdatesRequestV2 { since_metadata_version: state.last_synced },
187 )
188 .await?;
189
190 let tx = self.ro_tx().await;
191 let db = tx.db();
192
193 let mut without_orphans = Vec::new();
195 let me = Owner(self.keychain.get_pk()?);
196 let remote = db.base_metadata.stage(updates.file_metadata).to_lazy();
197 for id in remote.tree.staged.ids() {
198 let meta = remote.find(&id)?;
199 if remote.maybe_find_parent(meta).is_some()
200 || meta
201 .user_access_keys()
202 .iter()
203 .any(|k| k.encrypted_for == me.0)
204 {
205 without_orphans.push(remote.find(&id)?.clone());
206 }
207 }
208
209 let pruned_tree = db.base_metadata.stage(without_orphans).pruned()?.to_lazy();
211 let (_, deduped_changes) = pruned_tree.unstage();
212
213 let mut root_id = None;
215 if db.root.get().is_none() {
216 let root = deduped_changes
217 .all_files()?
218 .into_iter()
219 .find(|f| f.is_root())
220 .ok_or(LbErrKind::RootNonexistent)?;
221 root_id = Some(*root.id());
222 }
223
224 state.remote_changes = deduped_changes;
225 state.updates_as_of = updates.as_of_metadata_version;
226 state.new_root = root_id;
227
228 Ok(())
229 }
230
231 async fn fetch_required_docs(&self, state: &mut SyncState) -> LbResult<()> {
232 let mut docs_to_pull = vec![];
233
234 let tx = self.ro_tx().await;
235 let db = tx.db();
236
237 let mut files_with_local_edits = vec![];
238 let local = db.base_metadata.stage(&db.local_metadata);
239 for id in local.staged.ids() {
240 if let Some(base) = local.base.maybe_find(&id) {
241 if let Some(local_hmac) = local.find(&id)?.document_hmac() {
242 if Some(local_hmac) != base.document_hmac() {
243 files_with_local_edits.push(id);
244 println!("local edits found");
245 }
246 }
247 }
248 }
249
250 let mut remote = db
251 .base_metadata
252 .stage(state.remote_changes.clone())
253 .to_lazy();
254
255 for id in remote.tree.staged.ids() {
256 if remote.calculate_deleted(&id)? {
257 continue;
258 }
259 let remote_hmac = remote.find(&id)?.document_hmac().cloned();
260 let base_hmac = remote
261 .tree
262 .base
263 .maybe_find(&id)
264 .and_then(|f| f.document_hmac())
265 .cloned();
266 if base_hmac == remote_hmac {
267 continue;
268 }
269
270 if let Some(remote_hmac) = remote_hmac {
271 if self.docs.exists(id, base_hmac) && !self.docs.exists(id, Some(remote_hmac)) {
275 docs_to_pull.push((id, remote_hmac));
276 }
277
278 if files_with_local_edits.contains(&id)
281 && !docs_to_pull
282 .iter()
283 .any(|(already_pulling, _)| already_pulling == &id)
284 {
285 if let Some(base_hmac) = base_hmac {
286 if !self.docs.exists(id, Some(base_hmac)) {
287 docs_to_pull.push((id, base_hmac));
290 }
291 }
292 docs_to_pull.push((id, remote_hmac));
293 }
294 }
295 }
296 drop(tx);
297
298 let futures = docs_to_pull
299 .into_iter()
300 .map(|(id, hmac)| async move { self.ensure_doc_available(id, hmac).await.map(|_| id) });
301
302 let mut stream = stream::iter(futures).buffer_unordered(
303 thread::available_parallelism()
304 .unwrap_or(NonZeroUsize::new(4).unwrap())
305 .into(),
306 );
307
308 while let Some(fut) = stream.next().await {
309 let id = fut?;
310 state.pulled_docs.push(id);
311 }
312
313 Ok(())
314 }
315
316 pub(crate) async fn ensure_doc_available(
317 &self, id: Uuid, hmac: DocumentHmac,
318 ) -> LbResult<Option<EncryptedDocument>> {
319 if self.docs.exists(id, Some(hmac)) {
324 return Ok(None);
325 }
326
327 self.events
328 .sync_update(SyncIncrement::PullingDocument(id, true));
329 let remote_document = self
330 .client
331 .request(self.get_account()?, GetDocRequest { id, hmac })
332 .await?;
333 self.docs
334 .insert(id, Some(hmac), &remote_document.content)
335 .await?;
336 self.events
337 .sync_update(SyncIncrement::PullingDocument(id, false));
338
339 Ok(Some(remote_document.content))
340 }
341
342 pub(crate) async fn fetch_doc(
343 &self, id: Uuid, hmac: DocumentHmac,
344 ) -> LbResult<EncryptedDocument> {
345 match self.ensure_doc_available(id, hmac).await? {
346 Some(doc) => Ok(doc),
347 None => self.docs.get(id, Some(hmac)).await,
348 }
349 }
350
351 async fn merge(&self, state: &mut SyncState) -> LbResult<()> {
354 let mut tx = self.begin_tx().await;
355 let db = tx.db();
356 let start = Instant::now();
357
358 let remote_changes = &state.remote_changes;
359
360 let me = Owner(self.keychain.get_pk()?);
362
363 let merge_changes = {
365 let mut base = (&db.base_metadata).to_lazy();
367 let remote_unlazy = (&db.base_metadata).to_staged(remote_changes);
368 let mut remote = remote_unlazy.as_lazy();
369 let mut local = (&db.base_metadata).to_staged(&db.local_metadata).to_lazy();
370
371 let mut files_to_unmove: HashSet<Uuid> = HashSet::new();
373 let mut files_to_unshare: HashSet<Uuid> = HashSet::new();
374 let mut links_to_delete: HashSet<Uuid> = HashSet::new();
375 let mut rename_increments: HashMap<Uuid, usize> = HashMap::new();
376 let mut duplicate_file_ids: HashMap<Uuid, Uuid> = HashMap::new();
377
378 'merge_construction: loop {
379 let mut deletions = {
381 let mut deletions = remote_unlazy.stage(Vec::new()).to_lazy();
382
383 let mut deletion_creations = HashSet::new();
385 for id in db.local_metadata.ids() {
386 if remote.maybe_find(&id).is_none() && !links_to_delete.contains(&id) {
387 deletion_creations.insert(id);
388 }
389 }
390 'drain_creations: while !deletion_creations.is_empty() {
391 'choose_a_creation: for id in &deletion_creations {
392 let id = *id;
394 let local_file = local.find(&id)?.clone();
395 let result = deletions.create_unvalidated(
396 id,
397 symkey::generate_key(),
398 local_file.parent(),
399 &local.name(&id, &self.keychain)?,
400 local_file.file_type(),
401 &self.keychain,
402 );
403 match result {
404 Ok(_) => {
405 deletion_creations.remove(&id);
406 continue 'drain_creations;
407 }
408 Err(ref err) => match err.kind {
409 LbErrKind::FileParentNonexistent => {
410 continue 'choose_a_creation;
411 }
412 _ => {
413 result?;
414 }
415 },
416 }
417 }
418 return Err(LbErrKind::Unexpected(format!(
419 "sync failed to find a topomodelal order for file creations: {deletion_creations:?}"
420 ))
421 .into());
422 }
423
424 for id in db.local_metadata.ids() {
426 let local_file = local.find(&id)?.clone();
427 if let Some(base_file) = db.base_metadata.maybe_find(&id).cloned() {
428 if !local_file.explicitly_deleted()
429 && local_file.parent() != base_file.parent()
430 && !files_to_unmove.contains(&id)
431 {
432 deletions.move_unvalidated(
434 &id,
435 local_file.parent(),
436 &self.keychain,
437 )?;
438 }
439 }
440 }
441
442 for id in db.local_metadata.ids() {
444 let local_file = local.find(&id)?.clone();
445 if local_file.explicitly_deleted() {
446 deletions.delete_unvalidated(&id, &self.keychain)?;
448 }
449 }
450 deletions
451 };
452
453 let mut merge = {
455 let mut merge = remote_unlazy.stage(Vec::new()).to_lazy();
456
457 let mut creations = HashSet::new();
459 for id in db.local_metadata.ids() {
460 if deletions.maybe_find(&id).is_some()
461 && !deletions.calculate_deleted(&id)?
462 && remote.maybe_find(&id).is_none()
463 && !links_to_delete.contains(&id)
464 {
465 creations.insert(id);
466 }
467 }
468 'drain_creations: while !creations.is_empty() {
469 'choose_a_creation: for id in &creations {
470 let id = *id;
472 let local_file = local.find(&id)?.clone();
473 let result = merge.create_unvalidated(
474 id,
475 local.decrypt_key(&id, &self.keychain)?,
476 local_file.parent(),
477 &local.name(&id, &self.keychain)?,
478 local_file.file_type(),
479 &self.keychain,
480 );
481 match result {
482 Ok(_) => {
483 creations.remove(&id);
484 continue 'drain_creations;
485 }
486 Err(ref err) => match err.kind {
487 LbErrKind::FileParentNonexistent => {
488 continue 'choose_a_creation;
489 }
490 _ => {
491 result?;
492 }
493 },
494 }
495 }
496 return Err(LbErrKind::Unexpected(format!(
497 "sync failed to find a topomodelal order for file creations: {creations:?}"
498 ))
499 .into());
500 }
501
502 for id in db.local_metadata.ids() {
505 if deletions.maybe_find(&id).is_none()
507 || deletions.calculate_deleted(&id)?
508 || (remote.maybe_find(&id).is_some()
509 && remote.calculate_deleted(&id)?)
510 {
511 continue;
512 }
513
514 let local_file = local.find(&id)?.clone();
515 let local_name = local.name(&id, &self.keychain)?;
516 let maybe_base_file = base.maybe_find(&id).cloned();
517 let maybe_remote_file = remote.maybe_find(&id).cloned();
518 if let Some(ref base_file) = maybe_base_file {
519 let base_name = base.name(&id, &self.keychain)?;
520 let remote_file = remote.find(&id)?.clone();
521 let remote_name = remote.name(&id, &self.keychain)?;
522
523 if local_file.parent() != base_file.parent()
525 && remote_file.parent() == base_file.parent()
526 && !files_to_unmove.contains(&id)
527 {
528 merge.move_unvalidated(&id, local_file.parent(), &self.keychain)?;
529 }
530
531 if local_name != base_name && remote_name == base_name {
533 merge.rename_unvalidated(&id, &local_name, &self.keychain)?;
534 }
535 }
536
537 let mut remote_keys = HashMap::new();
539 if let Some(ref remote_file) = maybe_remote_file {
540 for key in remote_file.user_access_keys() {
541 remote_keys.insert(
542 (Owner(key.encrypted_by), Owner(key.encrypted_for)),
543 (key.mode, key.deleted),
544 );
545 }
546 }
547 for key in local_file.user_access_keys() {
548 let (by, for_) = (Owner(key.encrypted_by), Owner(key.encrypted_for));
549 if let Some(&(remote_mode, remote_deleted)) =
550 remote_keys.get(&(by, for_))
551 {
552 if key.mode > remote_mode || !key.deleted && remote_deleted {
554 let mode = match key.mode {
555 UserAccessMode::Read => ShareMode::Read,
556 UserAccessMode::Write => ShareMode::Write,
557 UserAccessMode::Owner => continue,
558 };
559 merge.add_share_unvalidated(id, for_, mode, &self.keychain)?;
560 }
561 if key.deleted && !remote_deleted {
563 merge.delete_share_unvalidated(
564 &id,
565 Some(for_.0),
566 &self.keychain,
567 )?;
568 }
569 } else {
570 let mode = match key.mode {
572 UserAccessMode::Read => ShareMode::Read,
573 UserAccessMode::Write => ShareMode::Write,
574 UserAccessMode::Owner => continue,
575 };
576 merge.add_share_unvalidated(id, for_, mode, &self.keychain)?;
577 }
578 }
579
580 if files_to_unshare.contains(&id) {
582 merge.delete_share_unvalidated(&id, None, &self.keychain)?;
583 }
584
585 if let Some(&rename_increment) = rename_increments.get(&id) {
587 let name = NameComponents::from(&local_name)
588 .generate_incremented(rename_increment)
589 .to_name();
590 merge.rename_unvalidated(&id, &name, &self.keychain)?;
591 }
592
593 let base_hmac = maybe_base_file.and_then(|f| f.document_hmac().cloned());
595 let remote_hmac =
596 maybe_remote_file.and_then(|f| f.document_hmac().cloned());
597 let local_hmac = local_file.document_hmac().cloned();
598 if merge.access_mode(me, &id)? >= Some(UserAccessMode::Write)
599 && local_hmac != base_hmac
600 {
601 if remote_hmac != base_hmac && remote_hmac != local_hmac {
602 let merge_name = merge.name(&id, &self.keychain)?;
604 let document_type =
605 DocumentType::from_file_name_using_extension(&merge_name);
606
607 let base_document =
610 self.read_document_helper(id, &mut base).await?;
611 let remote_document =
612 self.read_document_helper(id, &mut remote).await?;
613 let local_document =
614 self.read_document_helper(id, &mut local).await?;
615
616 match document_type {
617 DocumentType::Text => {
618 let base_document =
621 String::from_utf8_lossy(&base_document).to_string();
622 let remote_document =
623 String::from_utf8_lossy(&remote_document).to_string();
624 let local_document =
625 String::from_utf8_lossy(&local_document).to_string();
626 let merged_document =
627 text::buffer::Buffer::from(base_document.as_str())
628 .merge(local_document, remote_document);
629 let encrypted_document = merge
630 .update_document_unvalidated(
631 &id,
632 &merged_document.into_bytes(),
633 &self.keychain,
634 )?;
635 let hmac = merge.find(&id)?.document_hmac().copied();
636 self.docs.insert(id, hmac, &encrypted_document).await?;
637 }
638 DocumentType::Drawing => {
639 let base_document =
640 String::from_utf8_lossy(&base_document).to_string();
641 let remote_document =
642 String::from_utf8_lossy(&remote_document).to_string();
643 let local_document =
644 String::from_utf8_lossy(&local_document).to_string();
645
646 let base_buffer = svg::buffer::Buffer::new(&base_document);
647 let remote_buffer =
648 svg::buffer::Buffer::new(&remote_document);
649 let mut local_buffer =
650 svg::buffer::Buffer::new(&local_document);
651
652 for (_, el) in local_buffer.elements.iter_mut() {
653 if let Element::Path(path) = el {
654 path.data.apply_transform(u_transform_to_bezier(
655 &Transform::from(
656 local_buffer
657 .weak_viewport_settings
658 .master_transform,
659 ),
660 ));
661 }
662 }
663 svg::buffer::Buffer::reload(
664 &mut local_buffer.elements,
665 &mut local_buffer.weak_images,
666 &mut local_buffer.weak_path_pressures,
667 &mut local_buffer.weak_viewport_settings,
668 &base_buffer,
669 &remote_buffer,
670 );
671
672 let merged_document = local_buffer.serialize();
673 let encrypted_document = merge
674 .update_document_unvalidated(
675 &id,
676 &merged_document.into_bytes(),
677 &self.keychain,
678 )?;
679 let hmac = merge.find(&id)?.document_hmac().copied();
680 self.docs.insert(id, hmac, &encrypted_document).await?;
681 }
682 DocumentType::Chat => {
683 let merged_document = chat::Buffer::merge(
685 &base_document,
686 &local_document,
687 &remote_document,
688 );
689 let encrypted_document = merge
690 .update_document_unvalidated(
691 &id,
692 &merged_document,
693 &self.keychain,
694 )?;
695 let hmac = merge.find(&id)?.document_hmac().copied();
696 self.docs.insert(id, hmac, &encrypted_document).await?;
697 }
698 DocumentType::Other => {
699 let merge_parent = *merge.find(&id)?.parent();
701 let duplicate_id = if let Some(&duplicate_id) =
702 duplicate_file_ids.get(&id)
703 {
704 duplicate_id
705 } else {
706 let duplicate_id = Uuid::new_v4();
707 duplicate_file_ids.insert(id, duplicate_id);
708 rename_increments.insert(duplicate_id, 1);
709 duplicate_id
710 };
711
712 let mut merge_name = merge_name;
713 merge_name = NameComponents::from(&merge_name)
714 .generate_incremented(
715 rename_increments
716 .get(&duplicate_id)
717 .copied()
718 .unwrap_or_default(),
719 )
720 .to_name();
721
722 merge.create_unvalidated(
723 duplicate_id,
724 symkey::generate_key(),
725 &merge_parent,
726 &merge_name,
727 FileType::Document,
728 &self.keychain,
729 )?;
730 let encrypted_document = merge
731 .update_document_unvalidated(
732 &duplicate_id,
733 &local_document,
734 &self.keychain,
735 )?;
736 let duplicate_hmac =
737 merge.find(&duplicate_id)?.document_hmac().copied();
738 self.docs
739 .insert(
740 duplicate_id,
741 duplicate_hmac,
742 &encrypted_document,
743 )
744 .await?;
745 }
746 }
747 } else {
748 let local_file = local.find(&id)?;
749 merge.overwrite_document_hmac_unvalidated(
750 &id,
751 local_file.document_hmac().copied(),
752 local_file.doc_size(),
753 &self.keychain,
754 )?;
755 }
756 }
757 }
758
759 for id in db.local_metadata.ids() {
762 if db.base_metadata.maybe_find(&id).is_some()
763 && deletions.calculate_deleted(&id)?
764 && !merge.calculate_deleted(&id)?
765 {
766 merge.delete_unvalidated(&id, &self.keychain)?;
768 }
769 }
770 for &id in &links_to_delete {
771 if merge.maybe_find(&id).is_some() && !merge.calculate_deleted(&id)? {
773 merge.delete_unvalidated(&id, &self.keychain)?;
774 }
775 }
776
777 merge
778 };
779
780 for link in merge.ids() {
782 if !merge.calculate_deleted(&link)? {
783 if let FileType::Link { target } = merge.find(&link)?.file_type() {
784 if merge.maybe_find(&target).is_some()
785 && merge.calculate_deleted(&target)?
786 {
787 if links_to_delete.insert(link) {
789 continue 'merge_construction;
790 } else {
791 return Err(LbErrKind::Unexpected(format!(
792 "sync failed to resolve broken link (deletion): {link:?}"
793 ))
794 .into());
795 }
796 }
797 }
798 }
799 }
800
801 let validate_result = merge.validate(me);
802 match validate_result {
803 Ok(_) => {
805 let (_, merge_changes) = merge.unstage();
806 break merge_changes;
807 }
808 Err(ref err) => match err.kind {
809 LbErrKind::Validation(ref vf) => match vf {
810 ValidationFailure::Cycle(ids) => {
812 let mut progress = false;
814 for &id in ids {
815 if db.local_metadata.maybe_find(&id).is_some()
816 && files_to_unmove.insert(id)
817 {
818 progress = true;
819 }
820 }
821 if !progress {
822 return Err(LbErrKind::Unexpected(format!(
823 "sync failed to resolve cycle: {ids:?}"
824 ))
825 .into());
826 }
827 }
828 ValidationFailure::PathConflict(ids) => {
829 let mut progress = false;
831 for &id in ids {
832 if duplicate_file_ids.values().any(|&dup| dup == id) {
833 *rename_increments.entry(id).or_insert(0) += 1;
834 progress = true;
835 break;
836 }
837 }
838 if !progress {
839 for &id in ids {
840 if db.local_metadata.maybe_find(&id).is_some() {
841 *rename_increments.entry(id).or_insert(0) += 1;
842 progress = true;
843 break;
844 }
845 }
846 }
847 if !progress {
848 return Err(LbErrKind::Unexpected(format!(
849 "sync failed to resolve path conflict: {ids:?}"
850 ))
851 .into());
852 }
853 }
854 ValidationFailure::SharedLink { link, shared_ancestor } => {
855 let mut progress = false;
857 if let Some(base_shared_ancestor) = base.maybe_find(shared_ancestor)
858 {
859 if !base_shared_ancestor.is_shared()
860 && files_to_unshare.insert(*shared_ancestor)
861 {
862 progress = true;
863 }
864 }
865 if !progress && links_to_delete.insert(*link) {
866 progress = true;
867 }
868 if !progress {
869 return Err(LbErrKind::Unexpected(format!(
870 "sync failed to resolve shared link: link: {link:?}, shared_ancestor: {shared_ancestor:?}"
871 )).into());
872 }
873 }
874 ValidationFailure::DuplicateLink { target } => {
875 let mut progress = false;
877 if let Some(link) = local.linked_by(target)? {
878 if links_to_delete.insert(link) {
879 progress = true;
880 }
881 }
882 if !progress {
883 return Err(LbErrKind::Unexpected(format!(
884 "sync failed to resolve duplicate link: target: {target:?}"
885 ))
886 .into());
887 }
888 }
889 ValidationFailure::BrokenLink(link) => {
890 if !links_to_delete.insert(*link) {
892 return Err(LbErrKind::Unexpected(format!(
893 "sync failed to resolve broken link: {link:?}"
894 ))
895 .into());
896 }
897 }
898 ValidationFailure::OwnedLink(link) => {
899 let mut progress = false;
901 if let Some(remote_link) = remote.maybe_find(link) {
902 if let FileType::Link { target } = remote_link.file_type() {
903 let remote_target = remote.find(&target)?;
904 if remote_target.owner() != me
905 && files_to_unmove.insert(target)
906 {
907 progress = true;
908 }
909 }
910 }
911 if !progress && links_to_delete.insert(*link) {
912 progress = true;
913 }
914 if !progress {
915 return Err(LbErrKind::Unexpected(format!(
916 "sync failed to resolve owned link: {link:?}"
917 ))
918 .into());
919 }
920 }
921 ValidationFailure::Orphan(_)
923 | ValidationFailure::NonFolderWithChildren(_)
924 | ValidationFailure::FileWithDifferentOwnerParent(_)
925 | ValidationFailure::FileNameTooLong(_)
926 | ValidationFailure::DeletedFileUpdated(_)
927 | ValidationFailure::NonDecryptableFileName(_) => {
928 validate_result?;
929 }
930 },
931 _ => {
933 validate_result?;
934 }
935 },
936 }
937 }
938 };
939
940 (&mut db.base_metadata)
942 .to_staged(remote_changes.clone())
943 .to_lazy()
944 .promote()?;
945 db.local_metadata.clear()?;
946 (&mut db.local_metadata)
947 .to_staged(merge_changes)
948 .to_lazy()
949 .promote()?;
950
951 db.base_metadata.stage(&mut db.local_metadata).prune()?;
954
955 if start.elapsed() > web_time::Duration::from_millis(100) {
956 warn!("sync merge held lock for {:?}", start.elapsed());
957 }
958
959 Ok(())
960 }
961
962 async fn send_pull_events(&self, state: &mut SyncState) -> LbResult<()> {
963 if state.new_root.is_some() {
964 self.events.signed_in();
965 }
966
967 if !state.remote_changes.is_empty() {
968 self.events.meta_changed(Actor::Sync);
969
970 let owner = Owner(self.keychain.get_pk()?);
971 if state.remote_changes.iter().any(|f| f.owner() != owner) {
972 self.events.pending_shares_changed();
973 }
974 }
975
976 for &doc in &state.pulled_docs {
977 self.events.doc_written(doc, Actor::Sync);
978 }
979
980 Ok(())
981 }
982
983 async fn commit_last_synced(&self, state: &mut SyncState) -> LbResult<()> {
984 let mut tx = self.begin_tx().await;
985 let db = tx.db();
986 db.last_synced.insert(state.updates_as_of as i64)?;
987
988 if let Some(root) = state.new_root {
989 db.root.insert(root)?;
990 }
991
992 Ok(())
993 }
994
995 async fn populate_pk_cache(&self) -> LbResult<()> {
996 let mut missing_owners = HashSet::new();
998 {
999 let tx = self.ro_tx().await;
1000 let db = tx.db();
1001 for file in db.base_metadata.get().values() {
1002 for user_access_key in file.user_access_keys() {
1003 let enc_by = Owner(user_access_key.encrypted_by);
1004 let enc_for = Owner(user_access_key.encrypted_for);
1005
1006 if !db.pub_key_lookup.get().contains_key(&enc_by) {
1007 missing_owners.insert(enc_by);
1008 }
1009
1010 if !db.pub_key_lookup.get().contains_key(&enc_for) {
1011 missing_owners.insert(enc_for);
1012 }
1013 }
1014 }
1015 }
1016
1017 let mut new_owners = HashMap::new();
1018 {
1019 for owner in missing_owners {
1020 let username_result = self
1021 .client
1022 .request(self.get_account().unwrap(), GetUsernameRequest { key: owner.0 })
1023 .await;
1024 new_owners.insert(owner, username_result);
1025 }
1026 }
1027
1028 let mut tx = self.begin_tx().await;
1029 let db = tx.db();
1030
1031 let have_updates = !new_owners.is_empty();
1032 for (owner, username) in new_owners {
1033 let username = match username {
1034 Err(ApiError::Endpoint(GetUsernameError::UserNotFound)) => "<unknown>".to_string(),
1035 Ok(username) => username.username,
1036 _ => continue, };
1038
1039 db.pub_key_lookup.insert(owner, username).unwrap();
1040 }
1041
1042 if have_updates {
1043 self.events.meta_changed(Actor::Sync);
1044 }
1045
1046 Ok(())
1047 }
1048
1049 async fn push_meta(&self) -> LbResult<()> {
1051 let mut updates = vec![];
1052 let mut local_changes_no_digests = Vec::new();
1053
1054 let tx = self.ro_tx().await;
1055 let db = tx.db();
1056
1057 let local = db.base_metadata.stage(&db.local_metadata).to_lazy();
1059
1060 for id in local.tree.staged.ids() {
1061 let mut local_change = local.tree.staged.find(&id)?.timestamped_value.value.clone();
1062 let maybe_base_file = local.tree.base.maybe_find(&id);
1063
1064 local_change.set_hmac_and_size(
1066 maybe_base_file.and_then(|f| f.document_hmac().copied()),
1067 maybe_base_file.and_then(|f| *f.timestamped_value.value.doc_size()),
1068 );
1069 let local_change = local_change.sign(&self.keychain)?;
1070
1071 local_changes_no_digests.push(local_change.clone());
1072 let file_diff = FileDiff { old: maybe_base_file.cloned(), new: local_change };
1073 updates.push(file_diff);
1074 }
1075
1076 drop(tx);
1077
1078 if !updates.is_empty() {
1079 self.client
1080 .request(self.get_account()?, UpsertRequestV2 { updates: updates.clone() })
1081 .await?;
1082 }
1083
1084 let mut tx = self.begin_tx().await;
1085 let db = tx.db();
1086
1087 (&mut db.base_metadata)
1089 .to_lazy()
1090 .stage(local_changes_no_digests)
1091 .promote()?;
1092 db.base_metadata.stage(&mut db.local_metadata).prune()?;
1093
1094 tx.end();
1095
1096 Ok(())
1097 }
1098
1099 async fn push_docs(&self) -> LbResult<()> {
1104 let mut updates = vec![];
1105 let mut local_changes_digests_only = vec![];
1106
1107 let tx = self.ro_tx().await;
1108 let db = tx.db();
1109 let start = Instant::now();
1110
1111 let local = db.base_metadata.stage(&db.local_metadata).to_lazy();
1112
1113 for id in local.tree.staged.ids() {
1114 let base_file = local.tree.base.find(&id)?.clone();
1115
1116 let mut local_change = base_file.timestamped_value.value.clone();
1118 local_change.set_hmac_and_size(
1119 local.find(&id)?.document_hmac().copied(),
1120 *local.find(&id)?.timestamped_value.value.doc_size(),
1121 );
1122
1123 if base_file.document_hmac() == local_change.document_hmac()
1124 || local_change.document_hmac().is_none()
1125 {
1126 continue;
1127 }
1128
1129 let local_change = local_change.sign(&self.keychain)?;
1130
1131 updates.push(FileDiff { old: Some(base_file), new: local_change.clone() });
1132 local_changes_digests_only.push(local_change);
1133 self.events
1134 .sync_update(SyncIncrement::PushingDocument(id, true));
1135 }
1136
1137 drop(tx);
1138 if start.elapsed() > web_time::Duration::from_millis(100) {
1139 warn!("sync push_docs held lock for {:?}", start.elapsed());
1140 }
1141
1142 let futures = updates.clone().into_iter().map(|diff| self.push_doc(diff));
1143
1144 let mut stream = stream::iter(futures).buffer_unordered(
1145 thread::available_parallelism()
1146 .unwrap_or(NonZeroUsize::new(4).unwrap())
1147 .into(),
1148 );
1149
1150 let mut docs_without_errors = vec![];
1151 let mut last_error: Option<LbErr> = None;
1152
1153 while let Some(fut) = stream.next().await {
1154 match fut {
1155 Ok(id) => {
1156 docs_without_errors.push(id);
1157 self.events
1158 .sync_update(SyncIncrement::PushingDocument(id, false));
1159 }
1160 Err(err) => {
1161 last_error = Some(err);
1162 }
1163 }
1164 }
1165
1166 local_changes_digests_only.retain(|f| docs_without_errors.contains(f.id()));
1167
1168 let mut tx = self.begin_tx().await;
1169 let db = tx.db();
1170 (&mut db.base_metadata)
1172 .to_lazy()
1173 .stage(local_changes_digests_only)
1174 .promote()?;
1175
1176 db.base_metadata.stage(&mut db.local_metadata).prune()?;
1177
1178 tx.end();
1179
1180 if let Some(err) = last_error { Err(err) } else { Ok(()) }
1181 }
1182
1183 async fn push_doc(&self, diff: FileDiff<SignedMeta>) -> LbResult<Uuid> {
1184 let id = *diff.new.id();
1185 let hmac = diff.new.document_hmac();
1186 let local_document_change = self.docs.get(id, hmac.copied()).await?;
1187 self.client
1188 .request(
1189 self.get_account()?,
1190 ChangeDocRequestV2 { diff, new_content: local_document_change },
1191 )
1192 .await?;
1193
1194 Ok(id)
1195 }
1196
1197 #[cfg(not(target_family = "wasm"))]
1198 async fn send_debug_info(&self, account: Account) {
1199 use crate::service::debug;
1200
1201 let max_panic_time = debug::latest_panic_time(&self.config.writeable_path)
1202 .await
1203 .unwrap_or_else(|e| {
1204 warn!("could not enumerate panic files: {e:?}");
1205 None
1206 });
1207
1208 let last_sent = {
1209 let tx = self.ro_tx().await;
1210 tx.db().last_extracted_panic.get().copied()
1211 };
1212
1213 let should_send = match (last_sent, max_panic_time) {
1214 (None, _) => true,
1215 (Some(prev), Some(cur)) => cur > prev,
1216 (Some(_), None) => false,
1217 };
1218
1219 if !should_send {
1220 return;
1221 }
1222
1223 let debug_info = self
1224 .debug_info("none provided - sync".to_string(), false)
1225 .await
1226 .unwrap();
1227
1228 let new_marker = max_panic_time.unwrap_or(0);
1229 let bg_self = self.clone();
1230 let task = async move {
1231 match bg_self
1232 .client
1233 .request(&account, UpsertDebugInfoRequest { debug_info })
1234 .await
1235 {
1236 Ok(_) => {
1237 let mut tx = bg_self.begin_tx().await;
1238 if let Err(e) = tx.db().last_extracted_panic.insert(new_marker) {
1239 warn!("could not record last_extracted_panic: {e:?}");
1240 }
1241 tx.end();
1242 }
1243 Err(e) => warn!("send_debug_info failed: {e:?}"),
1244 }
1245 };
1246
1247 if self.config.background_work {
1248 tokio::spawn(task);
1249 } else {
1250 task.await;
1251 }
1252 }
1253
1254 async fn read_document_helper<T>(
1255 &self, id: Uuid, tree: &mut LazyTree<T>,
1256 ) -> LbResult<DecryptedDocument>
1257 where
1258 T: TreeLike<F = SignedMeta>,
1259 {
1260 let file = tree.find(&id)?;
1261 validate::is_document(file)?;
1262 let hmac = file.document_hmac().copied();
1263
1264 if tree.calculate_deleted(&id)? {
1265 return Err(LbErrKind::FileNonexistent.into());
1266 }
1267
1268 let doc = match hmac {
1269 Some(hmac) => {
1270 let doc = self.docs.get(id, Some(hmac)).await?;
1271 tree.decrypt_document(&id, &doc, &self.keychain)?
1272 }
1273 None => vec![],
1274 };
1275
1276 Ok(doc)
1277 }
1278
1279 #[doc(hidden)]
1281 pub async fn server_dirty_ids(&self) -> LbResult<Vec<Uuid>> {
1282 let mut state = self.syncer.lock().await;
1283 self.inital_sync_state(&mut state).await?;
1284 self.process_deletions().await?;
1285 self.fetch_meta(&mut state).await?;
1286
1287 let server_ids = state.remote_changes.iter().map(|f| *f.id()).collect();
1288
1289 Ok(server_ids)
1290 }
1291
1292 pub(crate) fn setup_syncer(&self) {
1293 if self.config.background_work {
1294 self.clone().local_change_worker();
1295 self.clone().periodic_sync_worker();
1296 self.clone().post_sync_worker();
1297 }
1298 }
1299
1300 fn local_change_worker(self) {
1301 #[cfg(not(target_family = "wasm"))]
1302 tokio::spawn(async move {
1303 let mut events = self.subscribe();
1304
1305 let sync_criteria = |e: Event| {
1306 matches!(
1307 e,
1308 Event::MetadataChanged(Actor::User) | Event::DocumentWritten(_, Actor::User)
1309 )
1310 };
1311
1312 loop {
1313 time::sleep(Duration::from_millis(500)).await;
1314 let mut should_sync = false;
1315
1316 loop {
1318 let event = events.try_recv();
1319 match event {
1320 Ok(event) => {
1321 if sync_criteria(event) {
1322 should_sync = true;
1323 }
1324 }
1325 Err(TryRecvError::Empty) => break,
1326 _ => {
1327 panic!(
1328 "unexpected broadcast receive error, returning local_change_worker"
1329 );
1330 }
1331 }
1332 }
1333
1334 if !should_sync {
1337 let event = events.recv().await.unwrap();
1338 if sync_criteria(event) {
1339 self.sync().await.map_unexpected().log_and_ignore();
1340 } else {
1341 continue;
1342 }
1343 }
1344 }
1345 });
1346 }
1347
1348 fn periodic_sync_worker(self) {
1349 #[cfg(not(target_family = "wasm"))]
1350 tokio::spawn(async move {
1351 loop {
1352 if self.user_active().await {
1353 tokio::time::sleep(Duration::from_secs(3)).await;
1354 } else {
1355 tokio::time::sleep(Duration::from_secs(5 * 60)).await;
1356 }
1357 self.sync().await.map_unexpected().log_and_ignore();
1358 }
1359 });
1360 }
1361
1362 async fn user_active(&self) -> bool {
1363 let last_seen = self.user_last_seen.read().await;
1364 last_seen.elapsed() < Duration::from_secs(3 * 60)
1365 }
1366
1367 fn post_sync_worker(self) {
1368 #[cfg(not(target_family = "wasm"))]
1369 tokio::spawn(async move {
1370 let mut events = self.subscribe();
1371
1372 loop {
1373 let event = events.recv().await.unwrap();
1374 if let Event::Sync(SyncIncrement::SyncFinished(_)) = event {
1375 self.fetcher().await.map_unexpected().log_and_ignore();
1376 self.populate_pk_cache()
1377 .await
1378 .map_unexpected()
1379 .log_and_ignore();
1380 };
1381 }
1382 });
1383 }
1384
1385 async fn fetcher(&self) -> LbResult<()> {
1386 let mut files_to_pull = vec![];
1387
1388 let tx = self.ro_tx().await;
1389 let db = tx.db();
1390
1391 let Some(root) = db.root.get() else {
1392 return Ok(());
1393 };
1394
1395 let mut tree = db.base_metadata.stage(None).to_lazy();
1397
1398 for id in tree.descendants_using_links(root)? {
1399 let file = tree.find(&id)?;
1400 let hmac = file.document_hmac().copied();
1401
1402 if !file.is_document() {
1404 continue;
1405 }
1406
1407 if tree.calculate_deleted(&id)? {
1409 continue;
1410 }
1411
1412 if self.docs.exists(id, hmac) {
1413 continue;
1414 }
1415
1416 let name = tree.name(&id, &self.keychain)?;
1418 if !name.ends_with(".md") && !name.ends_with(".svg") {
1419 continue;
1420 }
1421
1422 files_to_pull.push((id, hmac));
1423 }
1424
1425 drop(tx);
1426
1427 for (id, hmac) in files_to_pull {
1431 if let Some(hmac) = hmac {
1432 self.ensure_doc_available(id, hmac).await?;
1433 }
1434 }
1435
1436 Ok(())
1437 }
1438}