From 2bc9ef84b98899dea6572117583ce797e4c0648e Mon Sep 17 00:00:00 2001 From: KiriAky 107 Date: Tue, 8 Sep 2026 15:25:49 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=B7=BB=E5=8A=A0=E6=8C=81=E4=B9=85?= =?UTF-8?q?=E5=8C=96=20Rust=20Sync=20=E4=B8=8A=E4=BC=A0=E9=98=9F=E5=88=97?= =?UTF-8?q?=E4=B8=8E=20HTTP=20=E4=BC=A0=E8=BE=93?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/contracts/Sync-v1契约.md | 8 + .../OpenNexus生产化实施进度-2026-09-08.md | 2 + frontend/src-tauri/src/lib.rs | 3 + frontend/src-tauri/src/sync_client.rs | 365 ++++++++++++++++++ frontend/src-tauri/src/sync_state.rs | 283 ++++++++++++++ frontend/src-tauri/src/workspace.rs | 18 +- frontend/src-tauri/tests/sync_push.rs | 139 +++++++ server sync/tests/host_fixture.py | 40 ++ 8 files changed, 851 insertions(+), 7 deletions(-) create mode 100644 frontend/src-tauri/src/sync_client.rs create mode 100644 frontend/src-tauri/src/sync_state.rs create mode 100644 frontend/src-tauri/tests/sync_push.rs create mode 100644 server sync/tests/host_fixture.py diff --git a/docs/contracts/Sync-v1契约.md b/docs/contracts/Sync-v1契约.md index fb866b7..d1cf361 100644 --- a/docs/contracts/Sync-v1契约.md +++ b/docs/contracts/Sync-v1契约.md @@ -2,6 +2,14 @@ 状态:服务端协议原型已实现,生产运维与桌面客户端尚未验收。入口为 `server sync/sync_server`;运行时 `/openapi.json` 是字段约束来源。不能将 SQLite Fixture 结果称为 PostgreSQL / MinIO 验收。 +## Rust 上传客户端增量 + +Workspace schema 3 新增持久绑定、上传作业与远端 heads。首次向已验证为空的远端绑定后,队列保留每次本地操作;大正文转入按摘要命名的 spool,清除已物化 outbox 的正文副本。提交基线来自远端确认值,首次发送时冻结,超时重试不得重算。确认响应逐字段核对后,heads、作业与 outbox 同事务更新。解绑封存旧绑定与队列,重新绑定从当前文件快照生成新操作,旧回调不能修改新绑定。 + +Rust HTTP 客户端已实现握手、登录、空远端复核、1 MiB 分块上传、查询 offset 续传、complete 和 Revision 提交。默认 HTTPS,测试 HTTP 必须显式启用;不跟随重定向,令牌不进入 URL,响应有大小限制。真实本地 HTTP Fixture 已验证 20 次编辑形成 20 个正确远端基线、同一提交重复 100 次无重复 revision。此 Fixture 使用 SQLite/磁盘对象,不替代生产 PostgreSQL/MinIO 验收。 + +此增量尚未开放桌面 Sync capability:Stronghold 会话接入、拉取/inbox、冲突 UI、附件故障矩阵及数据分类仍在实施。 + ## 身份与数据边界 生产 CLI 只接受 `postgresql+psycopg`,账号通过 `python -m sync_server create-user` 交互初始化。无开放注册、默认密码或内置共享账号。设备由每次登录注册;Access Token 有效 900 秒,Refresh Token 30 天,数据库只保存摘要。刷新轮换使旧会话立即失效。注销删除当前会话;设备撤销使该设备所有会话与上传立即不可访问。社区与 Sync 身份完全独立。 diff --git a/docs/development/OpenNexus生产化实施进度-2026-09-08.md b/docs/development/OpenNexus生产化实施进度-2026-09-08.md index dc3045d..1e4f17f 100644 --- a/docs/development/OpenNexus生产化实施进度-2026-09-08.md +++ b/docs/development/OpenNexus生产化实施进度-2026-09-08.md @@ -6,6 +6,8 @@ ## 持续实施增量 +- Rust 同步上传队列与 HTTP 客户端已实现首批链路:持久绑定、spool、冻结远端基线、offset 查询续传及原子确认。20 次离线编辑/重启重试、解绑隔离测试通过;真实本地 Sync 服务验证 20 条 revision 与 100 次提交重放。Sync capability 仍保持 false,等待双向同步、会话与 UI 完成。 + - 桌面全文/向量投影按 Vault 隔离,搜索前经 Host 对账文件摘要;向量重建从 Host 读取正文并保留 file_id。任务和笔记关联使用同一 Vault 的持久库。新增测试覆盖同路径双 Vault 隔离、变更/删除刷新、稳定 ID 和任务跨 Vault 不可见;真实 Core 测试增加全文搜索与删除后的检索验证。 - 检索隔离增量最终后端全量 899 项通过,真实 Core/Vault 搜索集成和 Rust desktop 全目标 Clippy 通过。语义重建的稳定 ID 测试使用显式测试 Embedding,不作为真实模型质量或性能证据。 diff --git a/frontend/src-tauri/src/lib.rs b/frontend/src-tauri/src/lib.rs index b40d17b..4440c1d 100644 --- a/frontend/src-tauri/src/lib.rs +++ b/frontend/src-tauri/src/lib.rs @@ -8,5 +8,8 @@ pub mod request_lifecycle; mod runtime_compat; #[cfg(windows)] pub mod session_lock; +#[cfg(feature = "desktop")] +pub mod sync_client; +pub mod sync_state; pub mod workspace; pub mod workspace_broker; diff --git a/frontend/src-tauri/src/sync_client.rs b/frontend/src-tauri/src/sync_client.rs new file mode 100644 index 0000000..b88dcef --- /dev/null +++ b/frontend/src-tauri/src/sync_client.rs @@ -0,0 +1,365 @@ +//! Bounded Sync v1 transport. No redirects, no token-bearing URLs, no implicit retries. +use crate::{ + sync_state::{Binding, Job}, + workspace::Workspace, +}; +use reqwest::{Client, Method, Url}; +use serde::{Deserialize, Serialize}; +use serde_json::{json, Value}; +use sha2::{Digest, Sha256}; +use std::{ + io::{Read, Seek, SeekFrom}, + sync::{Arc, Mutex}, + time::Duration, +}; +use zeroize::Zeroizing; + +#[derive(Debug)] +pub struct SyncError { + pub code: String, + pub status: u16, + pub retry_after: Option, +} +impl SyncError { + fn new(code: &str) -> Self { + Self { + code: code.into(), + status: 0, + retry_after: None, + } + } +} +type Result = std::result::Result; +pub trait WorkspaceAccess: Send + Sync { + fn access( + &self, + action: impl FnOnce(&mut Workspace) -> crate::workspace::Result, + ) -> Result; +} +impl WorkspaceAccess for Arc> { + fn access( + &self, + action: impl FnOnce(&mut Workspace) -> crate::workspace::Result, + ) -> Result { + action(&mut *self.lock().map_err(|_| SyncError::new("HOST_BUSY"))?).map_err(Into::into) + } +} +impl WorkspaceAccess for Arc>> { + fn access( + &self, + action: impl FnOnce(&mut Workspace) -> crate::workspace::Result, + ) -> Result { + action( + self.lock() + .map_err(|_| SyncError::new("HOST_BUSY"))? + .as_mut() + .ok_or_else(|| SyncError::new("WORKSPACE_NOT_OPEN"))?, + ) + .map_err(Into::into) + } +} +impl From for SyncError { + fn from(value: crate::workspace::HostError) -> Self { + Self::new(&value.code) + } +} +impl From for SyncError { + fn from(_: std::io::Error) -> Self { + Self::new("SYNC_IO_FAILED") + } +} + +/// Persist only via the Stronghold Sync scope, never as an IPC response. +#[derive(Serialize, Deserialize)] +pub struct Session { + pub access_token: String, + pub refresh_token: String, + pub expires_in: u64, + pub device_id: String, +} +impl Drop for Session { + fn drop(&mut self) { + use zeroize::Zeroize; + self.access_token.zeroize(); + self.refresh_token.zeroize(); + } +} +pub struct SyncClient { + endpoint: Url, + client: Client, + token: Zeroizing, +} + +impl SyncClient { + pub fn new(endpoint: &str, token: Zeroizing, allow_test_http: bool) -> Result { + let mut url = Url::parse(endpoint).map_err(|_| SyncError::new("SYNC_ENDPOINT_INVALID"))?; + if (url.scheme() != "https" && !(allow_test_http && url.scheme() == "http")) + || url.host_str().is_none() + || !url.username().is_empty() + || url.password().is_some() + || url.query().is_some() + || url.fragment().is_some() + || !matches!(url.path(), "" | "/") + { + return Err(SyncError::new("SYNC_ENDPOINT_INVALID")); + } + url.set_path("/"); + let client = Client::builder() + .timeout(Duration::from_secs(30)) + .redirect(reqwest::redirect::Policy::none()) + .build() + .map_err(|_| SyncError::new("SYNC_CLIENT_FAILED"))?; + Ok(Self { + endpoint: url, + client, + token, + }) + } + pub async fn json(&self, method: Method, path: &str, body: Option) -> Result { + self.send(method, path, body.map(|v| v.to_string().into_bytes()), true) + .await + } + async fn send( + &self, + method: Method, + path: &str, + body: Option>, + is_json: bool, + ) -> Result { + if !path.starts_with("sync/v1/") || path.contains(['\\', '#']) || path.contains("..") { + return Err(SyncError::new("SYNC_PATH_INVALID")); + } + let url = self + .endpoint + .join(path) + .map_err(|_| SyncError::new("SYNC_PATH_INVALID"))?; + let mut request = self.client.request(method, url); + if !self.token.is_empty() { + request = request.bearer_auth(self.token.as_str()); + } + if let Some(body) = body { + request = request + .header( + "Content-Type", + if is_json { + "application/json" + } else { + "application/octet-stream" + }, + ) + .body(body); + } + let mut response = request + .send() + .await + .map_err(|_| SyncError::new("SYNC_NETWORK_ERROR"))?; + let status = response.status().as_u16(); + let retry_after = response + .headers() + .get("retry-after") + .and_then(|v| v.to_str().ok()) + .and_then(|v| v.parse().ok()); + let mut body = Vec::new(); + while let Some(chunk) = response + .chunk() + .await + .map_err(|_| SyncError::new("SYNC_NETWORK_ERROR"))? + { + if body.len() + chunk.len() > 4 * 1024 * 1024 { + return Err(SyncError::new("SYNC_RESPONSE_TOO_LARGE")); + } + body.extend_from_slice(&chunk); + } + let value: Value = if body.is_empty() { + Value::Null + } else { + serde_json::from_slice(&body).map_err(|_| SyncError::new("SYNC_RESPONSE_INVALID"))? + }; + if !(200..300).contains(&status) { + let code = value["error"]["code"] + .as_str() + .filter(|v| v.len() <= 80 && v.bytes().all(|b| b.is_ascii_uppercase() || b == b'_')) + .unwrap_or("SYNC_HTTP_ERROR"); + return Err(SyncError { + code: code.into(), + status, + retry_after, + }); + } + Ok(value) + } + pub async fn login( + &self, + username: &str, + password: Zeroizing, + device_name: &str, + ) -> Result { + let value = self.json(Method::POST, "sync/v1/auth/sessions", Some(json!({"username":username,"password":password.as_str(),"device_name":device_name}))).await?; + serde_json::from_value(value).map_err(|_| SyncError::new("SYNC_RESPONSE_INVALID")) + } + pub async fn handshake(&self) -> Result<()> { + let result = self + .json(Method::GET, "sync/v1/handshake?protocol=1", None) + .await?; + if result["protocol"] != 1 + || result["chunk_size"] != 1048576 + || result["max_object_size"] != 104857600 + { + return Err(SyncError::new("PROTOCOL_INCOMPATIBLE")); + } + Ok(()) + } + pub async fn verify_empty(&self, remote_vault: &str) -> Result<()> { + identifier(remote_vault)?; + let page = self + .json( + Method::GET, + &format!("sync/v1/vaults/{remote_vault}/changes?limit=1"), + None, + ) + .await?; + if page["boundary"] != 0 || page["items"].as_array().is_none_or(|v| !v.is_empty()) { + return Err(SyncError::new("SYNC_RECONCILIATION_REQUIRED")); + } + Ok(()) + } + pub async fn push_one( + &self, + workspace: &impl WorkspaceAccess, + binding: &Binding, + ) -> Result { + if Url::parse(&binding.endpoint) + .ok() + .is_none_or(|url| url != self.endpoint) + { + return Err(SyncError::new("SYNC_BINDING_CHANGED")); + } + identifier(&binding.remote_vault)?; + let job = workspace.access(|ws| { + ws.sync_capture(&binding.id)?; + ws.sync_next(&binding.id) + })?; + let Some(job) = job else { + return Ok(false); + }; + if job.state == "conflict" { + return Err(SyncError::new("REVISION_CONFLICT")); + } + if job.operation == "put" && job.base_revision.is_none() { + self.upload(workspace, binding, &job).await?; + } + let payload = workspace.access(|ws| ws.sync_commit_payload(&job))?; + let revision = self + .json( + Method::POST, + &format!("sync/v1/vaults/{}/revisions", binding.remote_vault), + Some(payload), + ) + .await?; + workspace.access(|ws| ws.sync_ack(&job, &revision))?; + Ok(true) + } + async fn upload( + &self, + workspace: &impl WorkspaceAccess, + binding: &Binding, + job: &Job, + ) -> Result<()> { + let path = workspace.access(|ws| ws.sync_spool(&job.hash))?; + let mut file = std::fs::File::open(path)?; + if file.metadata()?.len() != job.size as u64 { + return Err(SyncError::new("SYNC_SPOOL_CORRUPT")); + } + let mut hasher = Sha256::new(); + let mut buffer = vec![0u8; 1048576]; + loop { + let count = file.read(&mut buffer)?; + if count == 0 { + break; + } + hasher.update(&buffer[..count]); + } + if format!("{:x}", hasher.finalize()) != job.hash { + return Err(SyncError::new("SYNC_SPOOL_CORRUPT")); + } + let base = format!("sync/v1/vaults/{}/uploads", binding.remote_vault); + let mut upload_id = job.upload_id.clone(); + let mut offset = 0; + if let Some(id) = &upload_id { + identifier(id)?; + match self.json(Method::GET, &format!("{base}/{id}"), None).await { + Ok(status) => { + offset = status["offset"] + .as_u64() + .filter(|v| *v <= job.size as u64) + .ok_or_else(|| SyncError::new("SYNC_RESPONSE_INVALID"))? + } + Err(error) + if matches!(error.code.as_str(), "UPLOAD_EXPIRED" | "UPLOAD_DAMAGED") => + { + upload_id = None + } + Err(error) => return Err(error), + } + } + if upload_id.is_none() { + let response = self + .json( + Method::POST, + &base, + Some(json!({"content_hash":job.hash,"size":job.size})), + ) + .await?; + if response["complete"] == true { + return Ok(()); + } + upload_id = Some( + response["upload_id"] + .as_str() + .ok_or_else(|| SyncError::new("SYNC_RESPONSE_INVALID"))? + .into(), + ); + workspace.access(|ws| ws.sync_upload(job, upload_id.as_deref()))?; + } + let id = upload_id.ok_or_else(|| SyncError::new("SYNC_RESPONSE_INVALID"))?; + identifier(&id)?; + file.seek(SeekFrom::Start(offset))?; + while offset < job.size as u64 { + workspace.access(|ws| ws.check_binding(&binding.id))?; + let count = file.read(&mut buffer)?; + if count == 0 { + return Err(SyncError::new("SYNC_SPOOL_CORRUPT")); + } + let value = self + .send( + Method::PUT, + &format!("{base}/{id}?offset={offset}"), + Some(buffer[..count].to_vec()), + false, + ) + .await?; + if value["offset"].as_u64() != Some(offset + count as u64) { + return Err(SyncError::new("SYNC_RESPONSE_INVALID")); + } + offset += count as u64; + } + let value = self + .json(Method::POST, &format!("{base}/{id}/complete"), None) + .await?; + if value["complete"] != true || value["content_hash"] != job.hash { + return Err(SyncError::new("SYNC_RESPONSE_INVALID")); + } + Ok(()) + } +} +fn identifier(value: &str) -> Result<()> { + if value.is_empty() + || value.len() > 80 + || !value + .bytes() + .all(|b| b.is_ascii_alphanumeric() || b == b'-') + { + return Err(SyncError::new("SYNC_IDENTIFIER_INVALID")); + } + Ok(()) +} diff --git a/frontend/src-tauri/src/sync_state.rs b/frontend/src-tauri/src/sync_state.rs new file mode 100644 index 0000000..de98448 --- /dev/null +++ b/frontend/src-tauri/src/sync_state.rs @@ -0,0 +1,283 @@ +//! Durable queue state. Network code never invents a remote base from a local revision. +use crate::workspace::{hash, HostError, Result, Workspace}; +use rusqlite::{params, OptionalExtension}; +use serde::{Deserialize, Serialize}; +use std::{fs, io::Write, path::PathBuf}; +use uuid::Uuid; + +#[derive(Clone, Serialize, Deserialize)] +pub struct Binding { + pub id: String, + pub endpoint: String, + pub remote_vault: String, + pub account: String, + pub cursor: i64, +} +#[derive(Clone, Serialize, Deserialize)] +pub struct Job { + pub binding: String, + pub operation_id: String, + pub file_id: String, + pub path: String, + pub hash: String, + pub size: i64, + pub operation: String, + pub state: String, + pub base_revision: Option, + pub upload_id: Option, +} + +impl Workspace { + pub fn sync_binding(&self) -> Result> { + Ok(self.db.query_row("SELECT id,endpoint,remote_vault,account,cursor FROM sync_bindings WHERE state='active'", [], |r| { + Ok(Binding { id:r.get(0)?, endpoint:r.get(1)?, remote_vault:r.get(2)?, account:r.get(3)?, cursor:r.get(4)? }) + }).optional()?) + } + pub(crate) fn check_binding(&self, binding: &str) -> Result<()> { + if self.sync_binding()?.is_none_or(|b| b.id != binding) { + return Err(HostError::new("SYNC_BINDING_CHANGED")); + } + Ok(()) + } + /// Caller verifies an empty remote and obtains a reconciliation confirmation first. + pub fn sync_bind_empty( + &mut self, + endpoint: &str, + remote_vault: &str, + account: &str, + ) -> Result { + if self.sync_binding()?.is_some() { + return Err(HostError::new("SYNC_ALREADY_BOUND")); + } + let had_binding: bool = + self.db + .query_row("SELECT EXISTS(SELECT 1 FROM sync_bindings)", [], |r| { + r.get(0) + })?; + let entries = self.scan()?; + let id = Uuid::new_v4().to_string(); + // Rebinding explicitly starts from the current snapshot, never an old account's queue. + if had_binding { + self.db.execute( + "UPDATE outbox SET state='archived' WHERE state IN ('pending','queued')", + [], + )?; + } + for entry in entries.into_iter().filter(|e| !e.is_folder && !e.deleted) { + let queued: bool = self.db.query_row( + "SELECT EXISTS(SELECT 1 FROM outbox WHERE file_id=?1 AND state='pending')", + [&entry.file_id], + |r| r.get(0), + )?; + if !queued { + let content = fs::read(self.resolve(&entry.path)?)?; + self.write(&entry.path, &entry.hash, &content, "local")?; + } + } + self.db.execute( + "INSERT INTO sync_bindings VALUES (?1,?2,?3,?4,'active',0)", + params![id, endpoint, remote_vault, account], + )?; + self.sync_capture(&id)?; + self.sync_binding()? + .ok_or_else(|| HostError::new("DATABASE_ERROR")) + } + pub fn sync_unbind(&mut self, binding: &str) -> Result<()> { + self.check_binding(binding)?; + let tx = self.db.transaction()?; + tx.execute( + "UPDATE sync_bindings SET state='archived' WHERE id=?1", + [binding], + )?; + tx.execute( + "UPDATE outbox SET state='archived' WHERE state IN ('pending','queued')", + [], + )?; + tx.commit()?; + Ok(()) + } + pub fn sync_spool(&self, digest: &str) -> Result { + if digest.len() != 64 + || !digest + .bytes() + .all(|v| v.is_ascii_hexdigit() && !v.is_ascii_uppercase()) + { + return Err(HostError::new("SYNC_HASH_INVALID")); + } + self.resolve(&format!("attachments/{digest}"))?; // Enforce the platform's general path rules. + let root = self.root.join(".ainote/sync-spool"); + if root.exists() { + let meta = fs::symlink_metadata(&root)?; + if !meta.is_dir() || meta.file_type().is_symlink() { + return Err(HostError::new("UNSAFE_PATH")); + } + #[cfg(windows)] + { + use std::os::windows::fs::MetadataExt; + if meta.file_attributes() & 0x400 != 0 { + return Err(HostError::new("UNSAFE_PATH")); + } + } + } + fs::create_dir_all(&root)?; + Ok(root.join(digest)) + } + pub fn sync_capture(&mut self, binding: &str) -> Result<()> { + self.check_binding(binding)?; + loop { + let pending = self.db.query_row("SELECT operation_id,file_id,path,hash,operation,content FROM outbox WHERE state='pending' ORDER BY rowid LIMIT 1", [], |r| { + Ok((r.get::<_,String>(0)?,r.get::<_,String>(1)?,r.get::<_,String>(2)?,r.get::<_,String>(3)?,r.get::<_,String>(4)?,r.get::<_,Vec>(5)?)) + }).optional()?; + let Some((operation_id, file_id, path, digest, operation, content)) = pending else { + break; + }; + if operation == "put" { + if hash(&content) != digest { + return Err(HostError::new("SYNC_SPOOL_CORRUPT")); + } + let target = self.sync_spool(&digest)?; + if target.exists() { + if fs::symlink_metadata(&target)?.file_type().is_symlink() + || hash(&fs::read(&target)?) != digest + { + return Err(HostError::new("SYNC_SPOOL_CORRUPT")); + } + } else { + let mut temp = tempfile::NamedTempFile::new_in(target.parent().unwrap())?; + temp.write_all(&content)?; + temp.as_file().sync_all()?; + temp.persist_noclobber(target) + .map_err(|_| HostError::new("SYNC_SPOOL_FAILED"))?; + } + } + let tx = self.db.transaction()?; + tx.execute("INSERT OR IGNORE INTO sync_jobs VALUES (?1,?2,?3,?4,?5,?6,?7,'pending',NULL,NULL,NULL,NULL)", + params![binding,operation_id,file_id,path,digest,content.len() as i64,operation])?; + tx.execute( + "UPDATE outbox SET state='queued',content=X'' WHERE operation_id=?1", + [&operation_id], + )?; + tx.commit()?; + } + Ok(()) + } + pub fn sync_next(&self, binding: &str) -> Result> { + self.check_binding(binding)?; + Ok(self.db.query_row("SELECT binding,operation_id,file_id,path,hash,size,operation,state,base_revision,upload_id FROM sync_jobs WHERE binding=?1 AND state NOT IN ('acked','archived') ORDER BY rowid LIMIT 1", [binding], |r| { + Ok(Job { binding:r.get(0)?,operation_id:r.get(1)?,file_id:r.get(2)?,path:r.get(3)?,hash:r.get(4)?,size:r.get(5)?,operation:r.get(6)?,state:r.get(7)?,base_revision:r.get(8)?,upload_id:r.get(9)? }) + }).optional()?) + } + pub fn sync_upload(&self, job: &Job, upload: Option<&str>) -> Result<()> { + self.check_binding(&job.binding)?; + self.db.execute("UPDATE sync_jobs SET state='uploading',upload_id=?3 WHERE binding=?1 AND operation_id=?2 AND base_revision IS NULL", params![job.binding,job.operation_id,upload])?; + Ok(()) + } + pub fn sync_commit_payload(&self, job: &Job) -> Result { + self.check_binding(&job.binding)?; + // The base is frozen exactly once. A response loss reuses the byte-equivalent payload. + self.db.execute("UPDATE sync_jobs SET state='committing',base_revision=COALESCE((SELECT revision FROM sync_heads WHERE binding=?1 AND file_id=?3),0) WHERE binding=?1 AND operation_id=?2 AND base_revision IS NULL", + params![job.binding,job.operation_id,job.file_id])?; + let base: i64 = self.db.query_row( + "SELECT base_revision FROM sync_jobs WHERE binding=?1 AND operation_id=?2", + params![job.binding, job.operation_id], + |r| r.get(0), + )?; + Ok( + serde_json::json!({"operation_id":job.operation_id,"file_id":job.file_id,"base_revision":base,"path":job.path, + "operation":job.operation,"content_hash":if job.operation=="put" {Some(&job.hash)} else {None},"size":job.size}), + ) + } + pub fn sync_ack(&mut self, job: &Job, revision: &serde_json::Value) -> Result<()> { + self.check_binding(&job.binding)?; + let payload = self.sync_commit_payload(job)?; + for field in [ + "operation_id", + "file_id", + "base_revision", + "path", + "operation", + "size", + ] { + if revision[field] != payload[field] { + return Err(HostError::new("SYNC_RESPONSE_INVALID")); + } + } + if revision["hash"] != payload["content_hash"] + || revision["vault_id"] + != self + .sync_binding()? + .ok_or_else(|| HostError::new("SYNC_BINDING_CHANGED"))? + .remote_vault + { + return Err(HostError::new("SYNC_RESPONSE_INVALID")); + } + let sequence = revision["sequence"] + .as_i64() + .filter(|v| *v > payload["base_revision"].as_i64().unwrap_or(0)) + .ok_or_else(|| HostError::new("SYNC_RESPONSE_INVALID"))?; + let tx = self.db.transaction()?; + tx.execute("INSERT INTO sync_heads VALUES (?1,?2,?3,?4,?5) ON CONFLICT(binding,file_id) DO UPDATE SET revision=excluded.revision,path=excluded.path,hash=excluded.hash WHERE sync_heads.revision = std::result::Result; impl HostError { - fn new(code: &str) -> Self { + pub(crate) fn new(code: &str) -> Self { Self { code: code.into(), message: code.into(), @@ -100,7 +100,7 @@ pub fn portable_path_string(path: &Path) -> String { pub struct Workspace { pub root: PathBuf, pub vault_id: String, - db: Connection, + pub(crate) db: Connection, _lock: File, } @@ -135,12 +135,12 @@ impl Workspace { let db = Connection::open(db_path)?; db.execute_batch("PRAGMA journal_mode=WAL; PRAGMA synchronous=FULL;")?; let version: i64 = db.query_row("PRAGMA user_version", [], |r| r.get(0))?; - if version > 2 { + if version > 3 { return Err(HostError::new("SCHEMA_INCOMPATIBLE")); } - if version == 1 { + if (1..3).contains(&version) { // Independent, complete SQLite backup before the schema ownership change. - let backup = managed.join(format!("host-schema1-{}.sqlite3", Uuid::new_v4())); + let backup = managed.join(format!("host-schema{version}-{}.sqlite3", Uuid::new_v4())); db.execute("VACUUM INTO ?1", [backup.to_string_lossy().as_ref()])?; } db.execute_batch("BEGIN IMMEDIATE; @@ -150,7 +150,11 @@ impl Workspace { CREATE TABLE IF NOT EXISTS file_ops (id TEXT PRIMARY KEY,kind TEXT NOT NULL,path TEXT NOT NULL,destination TEXT NOT NULL,hash TEXT NOT NULL,content BLOB NOT NULL,state TEXT NOT NULL DEFAULT 'pending'); CREATE TABLE IF NOT EXISTS outbox (operation_id TEXT PRIMARY KEY,file_id TEXT NOT NULL,revision INTEGER NOT NULL,path TEXT NOT NULL,hash TEXT NOT NULL,operation TEXT NOT NULL,content BLOB NOT NULL,state TEXT NOT NULL DEFAULT 'pending'); CREATE TABLE IF NOT EXISTS operations (operation_id TEXT PRIMARY KEY,fingerprint TEXT NOT NULL,state TEXT NOT NULL,result TEXT); - PRAGMA user_version=2; COMMIT;")?; + CREATE TABLE IF NOT EXISTS sync_bindings (id TEXT PRIMARY KEY,endpoint TEXT NOT NULL,remote_vault TEXT NOT NULL,account TEXT NOT NULL,state TEXT NOT NULL,cursor INTEGER NOT NULL DEFAULT 0); + CREATE UNIQUE INDEX IF NOT EXISTS sync_active ON sync_bindings(state) WHERE state='active'; + CREATE TABLE IF NOT EXISTS sync_jobs (binding TEXT NOT NULL,operation_id TEXT NOT NULL,file_id TEXT NOT NULL,path TEXT NOT NULL,hash TEXT NOT NULL,size INTEGER NOT NULL,operation TEXT NOT NULL,state TEXT NOT NULL,base_revision INTEGER,upload_id TEXT,remote_revision INTEGER,error TEXT,PRIMARY KEY(binding,operation_id)); + CREATE TABLE IF NOT EXISTS sync_heads (binding TEXT NOT NULL,file_id TEXT NOT NULL,revision INTEGER NOT NULL,path TEXT NOT NULL,hash TEXT NOT NULL,PRIMARY KEY(binding,file_id)); + PRAGMA user_version=3; COMMIT;")?; let vault_id: String = db .query_row("SELECT id FROM identity", [], |r| r.get(0)) .optional()? @@ -533,7 +537,7 @@ impl Workspace { pub fn pending_count(&self) -> Result { Ok(self.db.query_row( - "SELECT COUNT(*) FROM outbox WHERE state='pending'", + "SELECT COUNT(*) FROM outbox WHERE state IN ('pending','queued')", [], |r| r.get(0), )?) diff --git a/frontend/src-tauri/tests/sync_push.rs b/frontend/src-tauri/tests/sync_push.rs new file mode 100644 index 0000000..20c12b4 --- /dev/null +++ b/frontend/src-tauri/tests/sync_push.rs @@ -0,0 +1,139 @@ +#![cfg(feature = "desktop")] +use notesagent_host::{sync_client::SyncClient, workspace::Workspace}; +use serde_json::{json, Value}; +use std::{ + io::{BufRead, BufReader}, + path::Path, + process::{Child, Command, Stdio}, + sync::{Arc, Mutex}, +}; +use zeroize::Zeroizing; +struct Server(Child); +impl Drop for Server { + fn drop(&mut self) { + self.0.stdin.take(); + let _ = self.0.kill(); + let _ = self.0.wait(); + } +} + +#[tokio::test] +async fn actual_service_accepts_ordered_push_and_repeat_commit_without_duplicates() { + let root = tempfile::tempdir().unwrap(); + std::fs::write(root.path().join(".opennexus-test"), b"fixture").unwrap(); + let service = Path::new(env!("CARGO_MANIFEST_DIR")) + .join("../../server sync") + .canonicalize() + .unwrap(); + let python = service.join(if cfg!(windows) { + ".venv/Scripts/python.exe" + } else { + ".venv/bin/python" + }); + let mut server = Server( + Command::new(python) + .args(["-m", "tests.host_fixture"]) + .arg(root.path()) + .current_dir(service) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::null()) + .spawn() + .unwrap(), + ); + let mut line = String::new(); + BufReader::new(server.0.stdout.take().unwrap()) + .read_line(&mut line) + .unwrap(); + let ready: Value = serde_json::from_str(&line).unwrap(); + let endpoint = format!("http://127.0.0.1:{}", ready["port"]); + let public = SyncClient::new(&endpoint, Zeroizing::new(String::new()), true).unwrap(); + public.handshake().await.unwrap(); + let session = public + .login( + "rust-fixture", + Zeroizing::new("controlled-fixture-password".into()), + "Rust integration", + ) + .await + .unwrap(); + let client = SyncClient::new( + &endpoint, + Zeroizing::new(session.access_token.clone()), + true, + ) + .unwrap(); + let vault = client + .json( + reqwest::Method::POST, + "sync/v1/vaults", + Some(json!({"name":"Rust test"})), + ) + .await + .unwrap(); + let remote = vault["vault_id"].as_str().unwrap(); + client.verify_empty(remote).await.unwrap(); + let local = tempfile::tempdir().unwrap(); + let workspace = Arc::new(Mutex::new(Workspace::open(local.path()).unwrap())); + let binding = { + let mut ws = workspace.lock().unwrap(); + let mut digest = String::new(); + for index in 0..20 { + digest = ws + .write( + "note.md", + &digest, + format!("fixture-{index}").as_bytes(), + "local", + ) + .unwrap() + .hash; + } + ws.sync_bind_empty(&endpoint, remote, "rust-fixture") + .unwrap() + }; + let first = workspace + .lock() + .unwrap() + .sync_next(&binding.id) + .unwrap() + .unwrap(); + assert!(client.push_one(&workspace, &binding).await.unwrap()); + let payload = workspace + .lock() + .unwrap() + .sync_commit_payload(&first) + .unwrap(); + for _ in 0..100 { + let replay = client + .json( + reqwest::Method::POST, + &format!("sync/v1/vaults/{remote}/revisions"), + Some(payload.clone()), + ) + .await + .unwrap(); + assert_eq!(replay["sequence"], 1); + } + for _ in 1..20 { + assert!(client.push_one(&workspace, &binding).await.unwrap()); + } + assert!(!client.push_one(&workspace, &binding).await.unwrap()); + let changes = client + .json( + reqwest::Method::GET, + &format!("sync/v1/vaults/{remote}/changes"), + None, + ) + .await + .unwrap(); + assert_eq!(changes["items"].as_array().unwrap().len(), 20); + for (index, revision) in changes["items"].as_array().unwrap().iter().enumerate() { + assert_eq!(revision["base_revision"], index as i64); + assert_eq!(revision["sequence"], index as i64 + 1); + } + assert_eq!( + client.verify_empty(remote).await.unwrap_err().code, + "SYNC_RECONCILIATION_REQUIRED" + ); +} diff --git a/server sync/tests/host_fixture.py b/server sync/tests/host_fixture.py new file mode 100644 index 0000000..5cdcf2e --- /dev/null +++ b/server sync/tests/host_fixture.py @@ -0,0 +1,40 @@ +"""Single-worker localhost fixture for real Rust HTTP interoperability, never deployment.""" +import asyncio +import json +from pathlib import Path +import socket +import sys +import threading +import uvicorn +from sync_server.app import create_app +from sync_server.database import Database +from sync_server.storage import DiskObjects + + +def main(): + root = Path(sys.argv[1]) + if not root.is_dir() or not (root / '.opennexus-test').is_file(): + raise SystemExit('ISOLATED_TEST_ROOT_REQUIRED') + database = Database('sqlite:///' + str(root / 'sync.sqlite3')) + app = create_app(database, DiskObjects(root / 'objects'), root / 'staging') + database.add_user('rust-fixture', 'controlled-fixture-password') + sock = socket.socket() + sock.bind(('127.0.0.1', 0)) + sock.listen(128) + server = uvicorn.Server(uvicorn.Config(app, log_config=None, access_log=False, timeout_graceful_shutdown=1)) + def parent(): + sys.stdin.buffer.read() + server.should_exit = True + threading.Thread(target=parent, daemon=True).start() + async def run(): + task = asyncio.create_task(server.serve(sockets=[sock])) + while not server.started: + if task.done(): await task; raise RuntimeError('FIXTURE_START_FAILED') + await asyncio.sleep(.01) + print(json.dumps({'port': sock.getsockname()[1]}), flush=True) + await task + try: asyncio.run(run()) + finally: sock.close(); database.engine.dispose() + + +if __name__ == '__main__': main()