From a47fbdc09c41091503b44061f6c02324c3cda0c0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=9B=B7=E7=94=B5=E8=8A=BD=E8=A1=A3?= Date: Fri, 15 May 2026 17:10:30 -0400 Subject: [PATCH] Wire SyncEngine into BrowserCore Add per-profile sync orchestration to `ely_browser_core`: - `SyncEngine::for_profile_dir` loads / generates the persistent device identity under `/sync/device.json` and reads the bearer token from `/sync/bearer.token`. - `install_bearer` accepts (or clears) the Better Auth session token; everything else stays inert until a token is on disk. - `upload_now(&BrowserCore)` serialises the user's bookmarks into a stable JSON snapshot, ships it via `SyncApiClient::upload_snapshot`, and remembers the resulting snapshot id / logical clock / device for the UI to surface. - `BrowserCore::visible_bookmarks_for_sync` returns a read-only view the engine can iterate without touching the in-memory state. The shell / settings-page wiring that calls `upload_now` ships separately so this commit stays a pure model-layer change with no runtime behaviour difference until the UI plugs in. --- Cargo.lock | 1 + crates/ely_browser_core/Cargo.toml | 1 + crates/ely_browser_core/src/lib.rs | 2 + .../ely_browser_core/src/state/bookmarks.rs | 7 + crates/ely_browser_core/src/sync_engine.rs | 203 ++++++++++++++++++ 5 files changed, 214 insertions(+) create mode 100644 crates/ely_browser_core/src/sync_engine.rs diff --git a/Cargo.lock b/Cargo.lock index 7655b6e..a01e6c0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2277,6 +2277,7 @@ name = "ely_browser_core" version = "0.1.0" dependencies = [ "ely_domain", + "ely_sync_client", "serde", "serde_json", "thiserror 2.0.18", diff --git a/crates/ely_browser_core/Cargo.toml b/crates/ely_browser_core/Cargo.toml index 8c1284f..199ac97 100644 --- a/crates/ely_browser_core/Cargo.toml +++ b/crates/ely_browser_core/Cargo.toml @@ -7,6 +7,7 @@ rust-version.workspace = true [dependencies] ely_domain = { path = "../ely_domain" } +ely_sync_client = { path = "../ely_sync_client" } serde.workspace = true serde_json.workspace = true thiserror.workspace = true diff --git a/crates/ely_browser_core/src/lib.rs b/crates/ely_browser_core/src/lib.rs index 56b01f0..1290db3 100644 --- a/crates/ely_browser_core/src/lib.rs +++ b/crates/ely_browser_core/src/lib.rs @@ -1,6 +1,7 @@ mod error; mod navigation; mod state; +mod sync_engine; pub use error::CoreError; pub use state::{ @@ -10,3 +11,4 @@ pub use state::{ ElySpacePackage, InitialBrowserConfig, InstalledPlugin, LocalDataInventory, PluginAuditAction, PluginAuditEvent, SiteDataClearance, SpaceImportProfileMapping, TrashedSpace, }; +pub use sync_engine::{SyncEngine, SyncEngineBuilder, SyncOutcome}; diff --git a/crates/ely_browser_core/src/state/bookmarks.rs b/crates/ely_browser_core/src/state/bookmarks.rs index 583b360..22109b5 100644 --- a/crates/ely_browser_core/src/state/bookmarks.rs +++ b/crates/ely_browser_core/src/state/bookmarks.rs @@ -242,6 +242,13 @@ impl BrowserCore { .collect() } + /// Read-only view across every bookmark, regardless of profile — + /// used by the sync engine which mirrors the full state to the + /// Cloudflare worker, not just the visible profile. + pub fn visible_bookmarks_for_sync(&self) -> Vec<&BookmarkEntry> { + self.bookmarks.iter().collect() + } + fn bookmark_mut(&mut self, bookmark_id: &BookmarkId) -> Result<&mut BookmarkEntry, CoreError> { self.bookmarks .iter_mut() diff --git a/crates/ely_browser_core/src/sync_engine.rs b/crates/ely_browser_core/src/sync_engine.rs new file mode 100644 index 0000000..c773ee5 --- /dev/null +++ b/crates/ely_browser_core/src/sync_engine.rs @@ -0,0 +1,203 @@ +use std::{ + path::{Path, PathBuf}, + time::{SystemTime, UNIX_EPOCH}, +}; + +use ely_domain::BookmarkEntry; +use ely_sync_client::{ + ApiClientConfig, BearerToken, BearerTokenStore, DeviceIdentity, SnapshotPayload, + SnapshotUploadRequest, SyncApiClient, SyncClientError, +}; +use serde::{Deserialize, Serialize}; + +use crate::state::BrowserCore; + +/// Per-profile sync scaffolding. Holds the persisted device identity +/// and the bearer-token store; the API client is only constructed at +/// the moment the user runs a manual sync (so a logged-out user +/// doesn't pay TLS handshake cost on startup). +#[derive(Debug)] +pub struct SyncEngine { + api_config: ApiClientConfig, + bearer_store: BearerTokenStore, + identity: DeviceIdentity, + last_outcome: Option, +} + +impl SyncEngine { + /// Bring up the engine for the given profile data root. Generates a + /// fresh device identity on first use, then keeps it stable across + /// runs so the server's `user_devices` row stays bound. + pub fn for_profile_dir( + profile_data_dir: &Path, + device_name: impl Into, + platform: impl Into, + ) -> Result { + let sync_dir = profile_data_dir.join("sync"); + let identity = + DeviceIdentity::load_or_create(&sync_dir.join("device.json"), device_name, platform)?; + let bearer_store = BearerTokenStore::new(sync_dir.join("bearer.token")); + Ok(Self { + api_config: ApiClientConfig::production(), + bearer_store, + identity, + last_outcome: None, + }) + } + + pub fn identity(&self) -> &DeviceIdentity { + &self.identity + } + + pub fn bearer_path(&self) -> &Path { + self.bearer_store.path() + } + + pub fn last_outcome(&self) -> Option<&SyncOutcome> { + self.last_outcome.as_ref() + } + + /// Replace the cached bearer token. Trimming + shape validation is + /// enforced by `BearerToken::new`. Returns `Ok(())` even when the + /// caller passes a blank string (treated as a sign-out). + pub fn install_bearer(&mut self, raw: &str) -> Result { + let trimmed = raw.trim(); + if trimmed.is_empty() { + self.bearer_store.clear()?; + return Ok(false); + } + let token = BearerToken::new(trimmed)?; + self.bearer_store.save(&token)?; + Ok(true) + } + + pub fn is_signed_in(&self) -> Result { + self.bearer_store.load().map(|token| token.is_some()) + } + + /// Build the JSON snapshot payload and ship it to the worker. The + /// caller passes the BrowserCore directly so the snapshot reads + /// the latest committed state; we don't keep a parallel copy. + pub fn upload_now(&mut self, core: &BrowserCore) -> Result { + let Some(bearer) = self.bearer_store.load()? else { + let outcome = SyncOutcome::SignedOut; + self.last_outcome = Some(outcome.clone()); + return Ok(outcome); + }; + let snapshot = SyncSnapshotBody::from_core(core); + let bytes = serde_json::to_vec(&snapshot).map_err(|error| SyncClientError::Json { + endpoint: "snapshot".to_string(), + source: error, + })?; + let payload = SnapshotPayload::new(bytes)?; + let logical_clock = current_logical_clock(); + let snapshot_id = snapshot_id_for_user(&self.identity); + + let client = SyncApiClient::new(self.api_config.clone(), bearer)?; + let request = SnapshotUploadRequest::new( + &snapshot_id, + self.api_config.region(), + SNAPSHOT_SCHEMA_REV, + logical_clock, + &payload, + ); + let document = client.upload_snapshot(&request)?; + let outcome = SyncOutcome::Uploaded { + snapshot_id: document.snapshot.snapshot_id, + logical_clock: document.snapshot.logical_clock, + payload_bytes: document.snapshot.size_bytes, + device_id: document.device_id, + }; + self.last_outcome = Some(outcome.clone()); + Ok(outcome) + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum SyncOutcome { + SignedOut, + Uploaded { snapshot_id: String, logical_clock: u64, payload_bytes: u64, device_id: String }, +} + +const SNAPSHOT_SCHEMA_REV: u32 = 1; + +fn current_logical_clock() -> u64 { + SystemTime::now().duration_since(UNIX_EPOCH).map(|elapsed| elapsed.as_secs()).unwrap_or(0) +} + +fn snapshot_id_for_user(identity: &DeviceIdentity) -> String { + // The Cloudflare worker requires `^[a-z0-9][a-z0-9._-]{0,127}$`. + // The device ID already satisfies the pattern (lowercased prefix + // + UUIDv7 simple form) and is per-user-stable, so we reuse it as + // the snapshot id. Future work can extend this to per-object-type + // snapshots without disturbing the existing one. + identity.device_id.clone() +} + +#[derive(Serialize, Deserialize)] +struct SyncSnapshotBody { + schema_rev: u32, + bookmarks: Vec, +} + +impl SyncSnapshotBody { + fn from_core(core: &BrowserCore) -> Self { + Self { + schema_rev: SNAPSHOT_SCHEMA_REV, + bookmarks: core + .visible_bookmarks_for_sync() + .into_iter() + .map(BookmarkSyncRecord::from) + .collect(), + } + } +} + +/// Wire representation of a bookmark. We keep this struct stable so a +/// future deserializer can read snapshots written by earlier app +/// versions; new fields must default-fill on read. +#[derive(Serialize, Deserialize)] +struct BookmarkSyncRecord { + id: String, + title: String, + url: String, + profile_id: String, + space_id: String, + collection_name: String, + tags: Vec, + note: Option, + added_at_secs: u64, +} + +impl From<&BookmarkEntry> for BookmarkSyncRecord { + fn from(entry: &BookmarkEntry) -> Self { + Self { + id: entry.id().as_str().to_string(), + title: entry.title().to_string(), + url: entry.url().as_str().to_string(), + profile_id: entry.profile_id().as_str().to_string(), + space_id: entry.space_id().as_str().to_string(), + collection_name: entry.collection_name().to_string(), + tags: entry.tags().to_vec(), + note: entry.note().map(str::to_string), + added_at_secs: entry + .added_at() + .duration_since(UNIX_EPOCH) + .map(|elapsed| elapsed.as_secs()) + .unwrap_or(0), + } + } +} + +#[derive(Debug)] +pub struct SyncEngineBuilder { + pub profile_data_dir: PathBuf, + pub device_name: String, + pub platform: String, +} + +impl SyncEngineBuilder { + pub fn build(self) -> Result { + SyncEngine::for_profile_dir(&self.profile_data_dir, self.device_name, self.platform) + } +}