Skip to main content

lb_rs/subscribers/
syncer.rs

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    /// the starting point for updates for this sync pass
49    last_synced: u64,
50
51    /// if our pull is successful, this is the timestamp we will commit
52    updates_as_of: u64,
53
54    /// changes we pulled from the server, post deduplication
55    remote_changes: Vec<SignedMeta>,
56
57    /// did we pull a root on this pass?
58    new_root: Option<Uuid>,
59
60    /// what docs did we pull as a result of this sync
61    pulled_docs: Vec<Uuid>,
62}
63
64// we are gonna have a fetch metadata fn which will get the docs that it needs to get, the ones
65// that match should_fetch
66//
67// should_fetch is going to be a tree fn that will return true if:
68//     is md or svg that descends from
69
70impl 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        // todo: should this inform a re-pull?
107        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        // this loop implicitly prunes remote orphans
194        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        // this is what actually performs the deduplication
210        let pruned_tree = db.base_metadata.stage(without_orphans).pruned()?.to_lazy();
211        let (_, deduped_changes) = pruned_tree.unstage();
212
213        // initialize root if this is the first pull on this device
214        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                // pull a file if we have a prior base, this is our heuristic -- do they have the
272                // ability to edit this file while we release the lock and are pulling all the
273                // files
274                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                // this clause captures documents which went from being new -> multiple parties
279                // having updates. We'll still need the updates
280                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                            // this scenario basically only comes up in tests
288                            // someone modifies a file directly without reading the prior version
289                            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        // todo: in a lot of cases there is a list of ids we're trying to get, it would be better
320        // if the caller managed the event updates, the status would be more meaningful for longer
321
322        // in this world, can we get stuck pushing a doc as well? Probably fine for now
323        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    /// Pulls remote changes and constructs a changeset Merge such that Stage<Stage<Stage<Base, Remote>, Local>, Merge> is valid.
352    /// Promotes Base to Stage<Base, Remote> and Local to Stage<Local, Merge>
353    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        // fetch document updates and local documents for merge
361        let me = Owner(self.keychain.get_pk()?);
362
363        // compute merge changes
364        let merge_changes = {
365            // assemble trees
366            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            // changeset constraints - these evolve as we try to assemble changes and encounter validation failures
372            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                // process just the edits which allow us to check deletions in the result
380                let mut deletions = {
381                    let mut deletions = remote_unlazy.stage(Vec::new()).to_lazy();
382
383                    // creations
384                    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                            // create
393                            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                    // moves (creations happen first in case a file is moved into a new folder)
425                    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                                // move
433                                deletions.move_unvalidated(
434                                    &id,
435                                    local_file.parent(),
436                                    &self.keychain,
437                                )?;
438                            }
439                        }
440                    }
441
442                    // deletions (moves happen first in case a file is moved into a deleted folder)
443                    for id in db.local_metadata.ids() {
444                        let local_file = local.find(&id)?.clone();
445                        if local_file.explicitly_deleted() {
446                            // delete
447                            deletions.delete_unvalidated(&id, &self.keychain)?;
448                        }
449                    }
450                    deletions
451                };
452
453                // process all edits, dropping non-deletion edits for files that will be implicitly deleted
454                let mut merge = {
455                    let mut merge = remote_unlazy.stage(Vec::new()).to_lazy();
456
457                    // creations and edits of created documents
458                    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                            // create
471                            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                    // moves, renames, edits, and shares
503                    // creations happen first in case a file is moved into a new folder
504                    for id in db.local_metadata.ids() {
505                        // skip files that are already deleted or will be deleted
506                        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                            // move
524                            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                            // rename
532                            if local_name != base_name && remote_name == base_name {
533                                merge.rename_unvalidated(&id, &local_name, &self.keychain)?;
534                            }
535                        }
536
537                        // share
538                        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                                // upgrade share
553                                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                                // delete share
562                                if key.deleted && !remote_deleted {
563                                    merge.delete_share_unvalidated(
564                                        &id,
565                                        Some(for_.0),
566                                        &self.keychain,
567                                    )?;
568                                }
569                            } else {
570                                // add share
571                                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                        // share deletion due to conflicts
581                        if files_to_unshare.contains(&id) {
582                            merge.delete_share_unvalidated(&id, None, &self.keychain)?;
583                        }
584
585                        // rename due to path conflict
586                        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                        // edit
594                        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                                // merge
603                                let merge_name = merge.name(&id, &self.keychain)?;
604                                let document_type =
605                                    DocumentType::from_file_name_using_extension(&merge_name);
606
607                                // todo these accesses are potentially problematic
608                                // maybe not if service/docs is the persion doing network io
609                                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                                        // 3-way merge
619                                        // todo: a couple more clones than necessary
620                                        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                                        // line-union of append-only JSONL turns
684                                        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                                        // duplicate file
700                                        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                    // deletes
760                    // moves happen first in case a file is moved into a deleted folder
761                    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                            // delete
767                            merge.delete_unvalidated(&id, &self.keychain)?;
768                        }
769                    }
770                    for &id in &links_to_delete {
771                        // delete
772                        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                // validate; handle failures by introducing changeset constraints
781                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                                // delete links to deleted files
788                                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                    // merge changeset is valid
804                    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                            // merge changeset has resolvable validation errors and needs modification
811                            ValidationFailure::Cycle(ids) => {
812                                // revert all local moves in the cycle
813                                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                                // pick one local id and generate a non-conflicting filename
830                                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                                // if ancestor is newly shared, delete share, otherwise delete link
856                                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                                // delete local link with this target
876                                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                                // delete local link with this target
891                                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                                // if target is newly owned, unmove target, otherwise delete link
900                                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                            // merge changeset has unexpected validation errors
922                            ValidationFailure::Orphan(_)
923                            | ValidationFailure::NonFolderWithChildren(_)
924                            | ValidationFailure::FileWithDifferentOwnerParent(_)
925                            | ValidationFailure::FileNameTooLong(_)
926                            | ValidationFailure::DeletedFileUpdated(_)
927                            | ValidationFailure::NonDecryptableFileName(_) => {
928                                validate_result?;
929                            }
930                        },
931                        // merge changeset has unexpected errors
932                        _ => {
933                            validate_result?;
934                        }
935                    },
936                }
937            }
938        };
939
940        // base = remote; local = merge
941        (&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        // todo who else calls this did they manage locks right?
952        // self.cleanup_local_metadata()?;
953        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        // todo: is this the move?
997        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, // todo: possibly add some logging here
1037            };
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    /// Updates remote and base metadata to local.
1050    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        // remote = local
1058        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            // change everything but document hmac and re-sign
1065            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        // base = local
1088        (&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    /// Updates remote and base files to local. Assumes metadata is already pushed for all new files.
1100    // todo: make this so that all document updates are attempted and we don't just return the
1101    // first error. Once an attempt is made we can return any or all errors, either would be an
1102    // improvement
1103    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            // change only document hmac and re-sign
1117            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        // base = local (metadata)
1171        (&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    /// for tests only
1280    #[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(_))
1309                        | Event::DocumentWritten(_, Actor::User(_))
1310                )
1311            };
1312
1313            loop {
1314                time::sleep(Duration::from_millis(500)).await;
1315                let mut should_sync = false;
1316
1317                // drain the current channel, so we don't sync for each keystroke if they pile up
1318                loop {
1319                    let event = events.try_recv();
1320                    match event {
1321                        Ok(event) => {
1322                            if sync_criteria(event) {
1323                                should_sync = true;
1324                            }
1325                        }
1326                        Err(TryRecvError::Empty) => break,
1327                        _ => {
1328                            panic!(
1329                                "unexpected broadcast receive error, returning local_change_worker"
1330                            );
1331                        }
1332                    }
1333                }
1334
1335                // empty channel + nothing interesting has happened, sit and wait for something
1336                // interesting
1337                if !should_sync {
1338                    let event = events.recv().await.unwrap();
1339                    if sync_criteria(event) {
1340                        self.sync().await.map_unexpected().log_and_ignore();
1341                    } else {
1342                        continue;
1343                    }
1344                }
1345            }
1346        });
1347    }
1348
1349    fn periodic_sync_worker(self) {
1350        #[cfg(not(target_family = "wasm"))]
1351        tokio::spawn(async move {
1352            loop {
1353                if self.user_active().await {
1354                    tokio::time::sleep(Duration::from_secs(3)).await;
1355                } else {
1356                    tokio::time::sleep(Duration::from_secs(5 * 60)).await;
1357                }
1358                self.sync().await.map_unexpected().log_and_ignore();
1359            }
1360        });
1361    }
1362
1363    async fn user_active(&self) -> bool {
1364        let last_seen = self.user_last_seen.read().await;
1365        last_seen.elapsed() < Duration::from_secs(3 * 60)
1366    }
1367
1368    fn post_sync_worker(self) {
1369        #[cfg(not(target_family = "wasm"))]
1370        tokio::spawn(async move {
1371            let mut events = self.subscribe();
1372
1373            loop {
1374                let event = events.recv().await.unwrap();
1375                if let Event::Sync(SyncIncrement::SyncFinished(_)) = event {
1376                    self.fetcher().await.map_unexpected().log_and_ignore();
1377                    self.populate_pk_cache()
1378                        .await
1379                        .map_unexpected()
1380                        .log_and_ignore();
1381                };
1382            }
1383        });
1384    }
1385
1386    async fn fetcher(&self) -> LbResult<()> {
1387        let mut files_to_pull = vec![];
1388
1389        let tx = self.ro_tx().await;
1390        let db = tx.db();
1391
1392        let Some(root) = db.root.get() else {
1393            return Ok(());
1394        };
1395
1396        // we can only fetch things we know the server knows about
1397        let mut tree = db.base_metadata.stage(None).to_lazy();
1398
1399        for id in tree.descendants_using_links(root)? {
1400            let file = tree.find(&id)?;
1401            let hmac = file.document_hmac().copied();
1402
1403            // skip non-documents
1404            if !file.is_document() {
1405                continue;
1406            }
1407
1408            // skip deleted files
1409            if tree.calculate_deleted(&id)? {
1410                continue;
1411            }
1412
1413            if self.docs.exists(id, hmac) {
1414                continue;
1415            }
1416
1417            // skip non-first-party files
1418            let name = tree.name(&id, &self.keychain)?;
1419            if !name.ends_with(".md") && !name.ends_with(".svg") {
1420                continue;
1421            }
1422
1423            files_to_pull.push((id, hmac));
1424        }
1425
1426        drop(tx);
1427
1428        // this could all be done in parallel, but for now going to not do it that way
1429        // benefits: less work, but also ensures that a file that needs to be fetched immediately
1430        // can be
1431        for (id, hmac) in files_to_pull {
1432            if let Some(hmac) = hmac {
1433                self.ensure_doc_available(id, hmac).await?;
1434            }
1435        }
1436
1437        Ok(())
1438    }
1439}