feat: 完成 OpenNexus 第三阶段核心功能与生产化基础 #45
@@ -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 身份完全独立。
|
||||
|
||||
@@ -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,不作为真实模型质量或性能证据。
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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<u64>,
|
||||
}
|
||||
impl SyncError {
|
||||
fn new(code: &str) -> Self {
|
||||
Self {
|
||||
code: code.into(),
|
||||
status: 0,
|
||||
retry_after: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
type Result<T> = std::result::Result<T, SyncError>;
|
||||
pub trait WorkspaceAccess: Send + Sync {
|
||||
fn access<T>(
|
||||
&self,
|
||||
action: impl FnOnce(&mut Workspace) -> crate::workspace::Result<T>,
|
||||
) -> Result<T>;
|
||||
}
|
||||
impl WorkspaceAccess for Arc<Mutex<Workspace>> {
|
||||
fn access<T>(
|
||||
&self,
|
||||
action: impl FnOnce(&mut Workspace) -> crate::workspace::Result<T>,
|
||||
) -> Result<T> {
|
||||
action(&mut *self.lock().map_err(|_| SyncError::new("HOST_BUSY"))?).map_err(Into::into)
|
||||
}
|
||||
}
|
||||
impl WorkspaceAccess for Arc<Mutex<Option<Workspace>>> {
|
||||
fn access<T>(
|
||||
&self,
|
||||
action: impl FnOnce(&mut Workspace) -> crate::workspace::Result<T>,
|
||||
) -> Result<T> {
|
||||
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<crate::workspace::HostError> for SyncError {
|
||||
fn from(value: crate::workspace::HostError) -> Self {
|
||||
Self::new(&value.code)
|
||||
}
|
||||
}
|
||||
impl From<std::io::Error> 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<String>,
|
||||
}
|
||||
|
||||
impl SyncClient {
|
||||
pub fn new(endpoint: &str, token: Zeroizing<String>, allow_test_http: bool) -> Result<Self> {
|
||||
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<Value>) -> Result<Value> {
|
||||
self.send(method, path, body.map(|v| v.to_string().into_bytes()), true)
|
||||
.await
|
||||
}
|
||||
async fn send(
|
||||
&self,
|
||||
method: Method,
|
||||
path: &str,
|
||||
body: Option<Vec<u8>>,
|
||||
is_json: bool,
|
||||
) -> Result<Value> {
|
||||
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<String>,
|
||||
device_name: &str,
|
||||
) -> Result<Session> {
|
||||
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<bool> {
|
||||
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(())
|
||||
}
|
||||
@@ -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<i64>,
|
||||
pub upload_id: Option<String>,
|
||||
}
|
||||
|
||||
impl Workspace {
|
||||
pub fn sync_binding(&self) -> Result<Option<Binding>> {
|
||||
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<Binding> {
|
||||
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<PathBuf> {
|
||||
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<u8>>(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<Option<Job>> {
|
||||
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<serde_json::Value> {
|
||||
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<excluded.revision", params![job.binding,job.file_id,sequence,job.path,job.hash])?;
|
||||
tx.execute("UPDATE sync_jobs SET state='acked',remote_revision=?3,error=NULL WHERE binding=?1 AND operation_id=?2", params![job.binding,job.operation_id,sequence])?;
|
||||
tx.execute(
|
||||
"UPDATE outbox SET state='acked' WHERE operation_id=?1",
|
||||
[&job.operation_id],
|
||||
)?;
|
||||
tx.commit()?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
#[test]
|
||||
fn queue_uses_remote_bases_and_keeps_retry_payload_across_restart() {
|
||||
let root = tempfile::tempdir().unwrap();
|
||||
let mut ws = Workspace::open(root.path()).unwrap();
|
||||
let mut digest = String::new();
|
||||
for index in 0..20 {
|
||||
digest = ws
|
||||
.write("a.md", &digest, index.to_string().as_bytes(), "local")
|
||||
.unwrap()
|
||||
.hash;
|
||||
}
|
||||
let binding = ws
|
||||
.sync_bind_empty("https://sync.example", "fixture-vault", "fixture-user")
|
||||
.unwrap();
|
||||
for sequence in 1..=20 {
|
||||
let job = ws.sync_next(&binding.id).unwrap().unwrap();
|
||||
let payload = ws.sync_commit_payload(&job).unwrap();
|
||||
assert_eq!(payload["base_revision"], sequence - 1);
|
||||
drop(ws);
|
||||
ws = Workspace::open(root.path()).unwrap();
|
||||
assert_eq!(ws.sync_commit_payload(&job).unwrap(), payload);
|
||||
let response = serde_json::json!({"operation_id":job.operation_id,"file_id":job.file_id,"base_revision":sequence-1,
|
||||
"path":"a.md","operation":"put","hash":hash((sequence-1).to_string().as_bytes()),"size":(sequence-1).to_string().len(),
|
||||
"sequence":sequence,"vault_id":"fixture-vault"});
|
||||
ws.sync_ack(&job, &response).unwrap();
|
||||
}
|
||||
assert!(ws.sync_next(&binding.id).unwrap().is_none());
|
||||
assert_eq!(ws.read("a.md").unwrap().content, "19");
|
||||
}
|
||||
#[test]
|
||||
fn unbind_archives_old_work_and_stale_completion_is_rejected() {
|
||||
let root = tempfile::tempdir().unwrap();
|
||||
let mut ws = Workspace::open(root.path()).unwrap();
|
||||
ws.write("a.md", "", b"safe", "local").unwrap();
|
||||
let first = ws
|
||||
.sync_bind_empty("https://one.example", "remote-one", "account-one")
|
||||
.unwrap();
|
||||
let old = ws.sync_next(&first.id).unwrap().unwrap();
|
||||
ws.sync_unbind(&first.id).unwrap();
|
||||
let second = ws
|
||||
.sync_bind_empty("https://two.example", "remote-two", "account-two")
|
||||
.unwrap();
|
||||
let new = ws.sync_next(&second.id).unwrap().unwrap();
|
||||
assert_ne!(old.operation_id, new.operation_id);
|
||||
assert_eq!(
|
||||
ws.sync_commit_payload(&old).unwrap_err().code,
|
||||
"SYNC_BINDING_CHANGED"
|
||||
);
|
||||
assert_eq!(ws.read("a.md").unwrap().content, "safe");
|
||||
}
|
||||
}
|
||||
@@ -18,7 +18,7 @@ pub struct HostError {
|
||||
pub type Result<T> = std::result::Result<T, HostError>;
|
||||
|
||||
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<i64> {
|
||||
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),
|
||||
)?)
|
||||
|
||||
@@ -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"
|
||||
);
|
||||
}
|
||||
@@ -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()
|
||||
Reference in New Issue
Block a user