From cf39b1a905664f745edbb82fcc7899117c9628a0 Mon Sep 17 00:00:00 2001 From: KiriAky 107 Date: Tue, 8 Sep 2026 16:52:48 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E9=80=9A=E8=BF=87=20Host=20=E6=97=A5?= =?UTF-8?q?=E5=BF=97=E5=90=8C=E6=AD=A5=E7=B1=BB=E5=9E=8B=E5=8C=96=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E8=AE=B0=E5=BD=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- backend/app/services/desktop_tasks.py | 106 +++++++++ backend/app/services/task_service.py | 20 ++ backend/tests/test_desktop_projection.py | 28 ++- docs/contracts/Sync-v1契约.md | 11 + .../OpenNexus生产化实施进度-2026-09-08.md | 10 + frontend/src-tauri/src/lib.rs | 1 + frontend/src-tauri/src/records.rs | 217 ++++++++++++++++++ frontend/src-tauri/src/sync_discovery.rs | 6 + frontend/src-tauri/src/sync_state.rs | 4 + frontend/src-tauri/src/workspace.rs | 15 +- frontend/src-tauri/src/workspace_broker.rs | 63 +++++ frontend/src-tauri/tests/core_workspace.rs | 118 +++++++++- frontend/src-tauri/tests/sync_push.rs | 61 +++++ 13 files changed, 644 insertions(+), 16 deletions(-) create mode 100644 backend/app/services/desktop_tasks.py create mode 100644 frontend/src-tauri/src/records.rs diff --git a/backend/app/services/desktop_tasks.py b/backend/app/services/desktop_tasks.py new file mode 100644 index 0000000..126f2a6 --- /dev/null +++ b/backend/app/services/desktop_tasks.py @@ -0,0 +1,106 @@ +"""Desktop Task records are committed by Host before returning to Core callers.""" +from __future__ import annotations +from datetime import datetime, timezone +import re +from uuid import uuid4, uuid5, NAMESPACE_URL +from app import host_bridge +from app.contracts import Task, TaskStatus +from app.database.db import connect_knowledge, transaction +from app.errors import ApiError +from app.services import desktop_notes + + +def _call(method, **params): return desktop_notes.call('records.' + method, **params) +def _ms(value): return None if value is None else int(value.timestamp() * 1000) +def _datetime(value): return None if value is None else datetime.fromtimestamp(value / 1000, timezone.utc) +def _record(task): + return {'schema': 1, 'kind': 'task', 'id': task.task_id, 'data': { + 'title': task.title, 'description': task.description, 'status': task.status.value, + 'note_id': task.note_id, 'due_at_ms': _ms(task.due_at), + 'created_at_ms': _ms(task.created_at), 'updated_at_ms': _ms(task.updated_at)}} +def _task(record): + data = record['data'] + return Task(task_id=record['id'], title=data['title'], description=data['description'], status=data['status'], + note_id=data['note_id'], due_at=_datetime(data['due_at_ms']), created_at=_datetime(data['created_at_ms']), updated_at=_datetime(data['updated_at_ms'])) +def _operation(): return host_bridge.operation_id.get() or str(uuid4()) +def _replay(operation, task_id=None, values=None, deleted=False): + previous = _call('operation', operation_id=operation) + if previous is None: return None + if previous.get('state') != 'committed' or previous.get('deleted') != deleted: + raise ApiError(409, 'OPERATION_PAYLOAD_CONFLICT', '该操作标识已用于其他修改。') + task = _task(previous['record']) + if task_id is not None and task.task_id != task_id: + raise ApiError(409, 'OPERATION_PAYLOAD_CONFLICT', '该操作标识已用于其他任务。') + for name, value in (values or {}).items(): + actual = getattr(task, name) + if isinstance(actual, datetime) and isinstance(value, datetime): + actual, value = _ms(actual), _ms(value) + if actual != value: raise ApiError(409, 'OPERATION_PAYLOAD_CONFLICT', '该操作标识的字段不一致。') + return task + +def _migrate(): + # Only the already scoped Vault database is eligible; unassigned legacy global data stays untouched. + conn = connect_knowledge() + try: + if conn.execute("SELECT value FROM index_meta WHERE key='tasks_host_owned_v1'").fetchone(): return + from app.services.task_service import _task_from_row + for row in conn.execute('SELECT * FROM tasks ORDER BY task_id').fetchall(): + task = _task_from_row(row) + if _call('get', id=task.task_id) is None: + operation = str(uuid5(NAMESPACE_URL, 'opennexus-task-migration:' + host_bridge.vault_id.get() + ':' + task.task_id)) + _call('write', record=_record(task), expected='', operation_id=operation) + with transaction(conn): + conn.execute("INSERT OR REPLACE INTO index_meta VALUES ('tasks_host_owned_v1','1')") + finally: conn.close() + +def _link(note_id): + if not note_id: return None + try: return desktop_notes.call('read', file_id=note_id)['file_id'] + except ApiError as error: + if error.code == 'FILE_NOT_FOUND': raise ApiError(404, 'RESOURCE_NOT_FOUND', 'note not found', {'note_id': note_id}) from None + raise + +def create(*, title, description='', note_id=None, due_at=None): + _migrate(); operation = _operation() + values = {'title': title, 'description': description, 'note_id': note_id, 'due_at': due_at} + replay = _replay(operation, values=values) + if replay is not None: return replay + now = datetime.now(timezone.utc) + task_id = 'task_' + uuid5(NAMESPACE_URL, 'opennexus-task:' + operation).hex + task = Task(task_id=task_id, title=title, description=description, note_id=_link(note_id), due_at=due_at, created_at=now, updated_at=now) + receipt = _call('write', record=_record(task), expected='', operation_id=operation) + return _task(receipt['record']) +def get(task_id): + if re.fullmatch(r'task_[0-9a-f]{32}', task_id) is None: return None + _migrate(); value = _call('get', id=task_id) + return _task(value['record']) if value is not None else None +def list_tasks(*, limit, offset): + _migrate(); result = _call('list', limit=1000, offset=0); records = list(result['items']) + while len(records) < result['total']: + page = _call('list', limit=1000, offset=len(records)) + if not page['items']: break + records.extend(page['items']) + tasks = sorted((_task(value['record']) for value in records), key=lambda value: (value.updated_at, value.task_id), reverse=True) + return tasks[offset:offset+limit], len(tasks) +def update(task_id, values): + _migrate(); operation = _operation(); values = dict(values) + for key in ['title', 'description', 'status']: + if values.get(key) is None: values.pop(key, None) + if not set(values) <= {'title','description','status','note_id','due_at'}: raise ApiError(422, 'INVALID_ARGUMENT', '未知任务字段。') + replay = _replay(operation, task_id, values) + if replay is not None: return replay + current = _call('get', id=task_id) + if current is None: raise ApiError(404, 'RESOURCE_NOT_FOUND', 'task not found', {'task_id': task_id}) + if 'note_id' in values: values['note_id'] = _link(values['note_id']) + task = _task(current['record']).model_copy(update={**values, 'updated_at': datetime.now(timezone.utc)}) + if isinstance(task.status, str): task.status = TaskStatus(task.status) + receipt = _call('write', record=_record(task), expected=current['hash'], operation_id=operation) + return _task(receipt['record']) +def delete(task_id): + if re.fullmatch(r'task_[0-9a-f]{32}', task_id) is None: return False + _migrate(); operation = _operation() + if _replay(operation, task_id, deleted=True) is not None: return True + current = _call('get', id=task_id) + if current is None: return False + _call('delete', id=task_id, expected=current['hash'], operation_id=operation) + return True diff --git a/backend/app/services/task_service.py b/backend/app/services/task_service.py index c981274..c24ffff 100644 --- a/backend/app/services/task_service.py +++ b/backend/app/services/task_service.py @@ -13,6 +13,11 @@ from app.database.db import connect_knowledge as connect, transaction from app.errors import ApiError from app.operation_logs import log_event +def _desktop(): + from app.config import get_settings + return get_settings().environment == 'desktop' + + _write_locks = WeakKeyDictionary() @@ -63,6 +68,9 @@ def create_task( *, title: str, description: str = "", note_id: str | None = None, due_at: datetime | None = None, ) -> Task: + if _desktop(): + from app.services import desktop_tasks + return desktop_tasks.create(title=title, description=description, note_id=note_id, due_at=due_at) _prepare_note_link(note_id) if note_id and repository.get_note_record(note_id) is None: raise ApiError(404, "RESOURCE_NOT_FOUND", "note not found", {"note_id": note_id}) @@ -90,6 +98,9 @@ def create_task( def get_task(task_id: str) -> Task | None: + if _desktop(): + from app.services import desktop_tasks + return desktop_tasks.get(task_id) conn = connect() try: row = conn.execute("SELECT * FROM tasks WHERE task_id = ?", (task_id,)).fetchone() @@ -99,6 +110,9 @@ def get_task(task_id: str) -> Task | None: def list_tasks(*, limit: int, offset: int) -> tuple[list[Task], int]: + if _desktop(): + from app.services import desktop_tasks + return desktop_tasks.list_tasks(limit=limit, offset=offset) conn = connect() try: total = conn.execute("SELECT COUNT(*) FROM tasks").fetchone()[0] @@ -112,6 +126,9 @@ def list_tasks(*, limit: int, offset: int) -> tuple[list[Task], int]: def update_task(task_id: str, values: dict[str, object]) -> Task: + if _desktop(): + from app.services import desktop_tasks + return desktop_tasks.update(task_id, values) current = get_task(task_id) if current is None: raise ApiError(404, "RESOURCE_NOT_FOUND", "task not found", {"task_id": task_id}) @@ -154,6 +171,9 @@ def update_task(task_id: str, values: dict[str, object]) -> Task: def delete_task(task_id: str) -> bool: + if _desktop(): + from app.services import desktop_tasks + return desktop_tasks.delete(task_id) conn = connect() try: with transaction(conn): diff --git a/backend/tests/test_desktop_projection.py b/backend/tests/test_desktop_projection.py index 8650b29..6eb0904 100644 --- a/backend/tests/test_desktop_projection.py +++ b/backend/tests/test_desktop_projection.py @@ -40,18 +40,19 @@ def test_projection_isolates_same_path_and_refreshes_changed_deleted_content(tmp assert 'file-a' not in search(first, 'replacementtoken').model_dump_json() assert (tmp_path/'vault-state'/first/'core.sqlite3').is_file() assert (tmp_path/'vault-state'/second/'core.sqlite3').is_file() - from app.services import task_service token = host_bridge.vault_id.set(second) try: - task = task_service.create_task(title='Scoped task', note_id='file-b') - assert task.note_id == 'file-b' - finally: - host_bridge.vault_id.reset(token) + conn = db.connect_knowledge() + with db.transaction(conn): + conn.execute("INSERT INTO tasks VALUES ('legacy-task','Scoped task','','todo','file-b',NULL,'2026-09-08T00:00:00+00:00','2026-09-08T00:00:00+00:00')") + conn.close() + finally: host_bridge.vault_id.reset(token) token = host_bridge.vault_id.set(first) try: - assert task_service.get_task(task.task_id) is None - finally: - host_bridge.vault_id.reset(token) + conn = db.connect_knowledge() + assert conn.execute("SELECT * FROM tasks WHERE task_id='legacy-task'").fetchone() is None + conn.close() + finally: host_bridge.vault_id.reset(token) def test_desktop_semantic_rebuild_preserves_host_file_id(tmp_path, monkeypatch): @@ -82,14 +83,19 @@ def test_host_identity_adoption_keeps_existing_task_links(tmp_path, monkeypatch) monkeypatch.setattr('app.config.get_settings', lambda: settings) document = {'file_id': 'before-merge', 'path': 'same.md', 'hash': sha256(b'test').hexdigest(), 'content': 'test', 'created_at': 0, 'updated_at': 1} monkeypatch.setattr(desktop_notes, 'call', lambda method, **params: {'items': [document], 'total': 1} if method == 'list' else document) - from app.services import task_service token = host_bridge.vault_id.set(str(uuid4())) try: - task = task_service.create_task(title='Preserve link', note_id='before-merge') + asyncio.run(desktop_projection.refresh()) + conn = db.connect_knowledge() + with db.transaction(conn): + conn.execute("INSERT INTO tasks VALUES ('legacy-task','Preserve link','','todo','before-merge',NULL,'2026-09-08T00:00:00+00:00','2026-09-08T00:00:00+00:00')") + conn.close() document['file_id'] = 'after-merge' document['aliases'] = ['before-merge'] asyncio.run(desktop_projection.refresh()) - assert task_service.get_task(task.task_id).note_id == 'after-merge' + conn = db.connect_knowledge() + assert conn.execute("SELECT note_id FROM tasks WHERE task_id='legacy-task'").fetchone()['note_id'] == 'after-merge' + conn.close() assert repository.get_note_record('before-merge') is None assert repository.get_note_record('after-merge') is not None finally: diff --git a/docs/contracts/Sync-v1契约.md b/docs/contracts/Sync-v1契约.md index 89ee4f0..eef36f5 100644 --- a/docs/contracts/Sync-v1契约.md +++ b/docs/contracts/Sync-v1契约.md @@ -53,3 +53,14 @@ Vault 行锁内执行幂等键验证、CAS、路径检查、序列分配及历 当前历史和 tombstone **永久保留**,游标不主动过期;不启用 GC,以免缺少恢复演练时删除历史对象。配额包含全部历史对象,未完成上传预留额度。通知尚未实现,客户端必须主动按 cursor 拉取。上传过期清理、对象 GC、限速指标、完整备份恢复工具与生产负载验收仍未交付。`/ready` 目前只检查数据库 schema,不证明 S3 就绪。 验证向量见 `server sync/tests/test_protocol.py`:双设备 CAS、幂等重放、移动与删除、历史恢复、稳定分页、隔离、撤销、刷新、配额、路径冲突、摘要与断点。 + + +## 逻辑任务记录 v1 + +普通文件传输增加保留命名空间 `opennexus-records/v1/tasks/task_<32位小写十六进制>.json`。文件 Revision 的 file_id 与任务业务 id 分开;记录仅包含 `schema=1`、`kind=task`、`id`、`data`。data 白名单为 title、description、status、note_id、due_at_ms、created_at_ms、updated_at_ms。时间采用 UTC Unix 毫秒,due_at_ms/note_id 可为空;status 为 todo/in_progress/done/cancelled。任务状态只表示任务状态,不启动目标设备上的 Agent 或后台作业。 + +Host 在写入 journal、捕获外部修改和上传前验证记录;未知字段、未知 schema/kind、非法路径或时间拒绝。记录不允许 api_key、token、environment、permissions 等附加字段。标题最多 4096 字节、说明最多 256 KiB、记录最多 1 MiB,Core list 响应最多 4 MiB。用户写入标题/说明的正文仍是用户内容,不按关键字审查正文。 + +桌面 Task CRUD 经 Host records broker,操作重放返回原始记录且同 ID 不同字段拒绝;Core 不以全局 SQLite 作为任务来源。已明确归属于当前 Vault 的旧 Task 表在首次访问时逐条迁移,完成后设置所有权标记,来源表保留;未分配 Vault 的全局旧数据不猜测归属。文件身份采纳时,任务链接通过 Host 别名解析,并在上传队列物化前将规范化引用写成新的逻辑记录。 + +本节仅交付任务逻辑记录。用户 Skill/配置、主题、可选对话等仍须逐类定义白名单与适配器;不得用复制任意 JSON/SQLite 代替。 diff --git a/docs/development/OpenNexus生产化实施进度-2026-09-08.md b/docs/development/OpenNexus生产化实施进度-2026-09-08.md index 63a8f3a..bfc7f62 100644 --- a/docs/development/OpenNexus生产化实施进度-2026-09-08.md +++ b/docs/development/OpenNexus生产化实施进度-2026-09-08.md @@ -90,3 +90,13 @@ Core 的独立数据目录目前不等于已授权 Vault。Python 旧笔记写 - 目前快照最多处理 100000 条历史 revision,超限明确拒绝;大库快照/扫描性能、任务和配置逻辑同步、完整分类故障矩阵及签名/沙箱等其他生产门槛继续实施。 - 本增量后端全量 900 项、Rust desktop 全目标 49 项通过(另 1 个父测试使用的进程辅助入口);Sync 界面 3 项交互测试、两个 TypeScript 项目检查和 Clippy `-D warnings` 通过。 + + +## 增量:Task 逻辑记录 + +- Task v1 使用专用命名空间和 Rust 字段白名单,秘密字段、未知字段/版本在 journal 前拒绝;普通配置文件仍不纳入同步。任务不依赖 Core SQLite 的全量复制,也不触发目标设备自动执行任务。 +- 桌面 Task CRUD 已交由 Host 持久文件和 outbox;同操作创建重复 20 次返回同一结果,不同字段重用 ID 拒绝;更新/删除同样通过操作回执。已有当前 Vault 任务逐条迁移,保留原表且可重复完成迁移。 +- 真实 Core HTTP 验证任务创建/更新/删除/重放、跨 Vault 拒绝、旧任务来源保留;真实 Sync HTTP 验证两客户端任务状态与删除传播。记录列表有 4 MiB 分页上限,读取有 1 MiB 文件上限。 +- 本地文件 ID 采纳后的 task note_id 通过别名解析,上传前生成持久规范化引用,避免把本机旧 ID 留给另一设备。其余逻辑记录类别继续实施。 + +- Task 增量后端全量 900 项、Rust desktop 全目标 50 项通过(另 1 个父测试驱动的进程辅助入口),Clippy `-D warnings` 通过。 diff --git a/frontend/src-tauri/src/lib.rs b/frontend/src-tauri/src/lib.rs index a4698b9..c3a4698 100644 --- a/frontend/src-tauri/src/lib.rs +++ b/frontend/src-tauri/src/lib.rs @@ -4,6 +4,7 @@ pub mod core; pub mod credentials; mod payloads; pub mod recent; +pub mod records; #[cfg(feature = "desktop")] pub mod request_lifecycle; mod runtime_compat; diff --git a/frontend/src-tauri/src/records.rs b/frontend/src-tauri/src/records.rs new file mode 100644 index 0000000..39c2e2a --- /dev/null +++ b/frontend/src-tauri/src/records.rs @@ -0,0 +1,217 @@ +//! Versioned logical records: explicit fields only, never raw application databases/config. +use crate::workspace::{HostError, Result, Workspace}; +use serde::{Deserialize, Serialize}; +use serde_json::{json, Value}; +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct TaskData { + pub title: String, + pub description: String, + pub status: String, + pub note_id: Option, + pub due_at_ms: Option, + pub created_at_ms: i64, + pub updated_at_ms: i64, +} +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Record { + pub schema: u32, + pub kind: String, + pub id: String, + pub data: TaskData, +} +pub fn path(id: &str) -> Result { + if !id.starts_with("task_") + || id.len() != 37 + || !id[5..] + .bytes() + .all(|v| v.is_ascii_digit() || (b'a'..=b'f').contains(&v)) + { + return Err(HostError::new("RECORD_ID_INVALID")); + } + Ok(format!("opennexus-records/v1/tasks/{id}.json")) +} +pub fn is_record(path: &str) -> bool { + path.starts_with("opennexus-records/") +} +pub fn allowed(path_value: &str) -> bool { + path_value + .strip_prefix("opennexus-records/v1/tasks/") + .and_then(|v| v.strip_suffix(".json")) + .is_some_and(|id| path(id).is_ok()) +} +pub fn validate(path_value: &str, content: &[u8]) -> Result { + if content.len() > 1024 * 1024 { + return Err(HostError::new("RECORD_TOO_LARGE")); + } + let record: Record = + serde_json::from_slice(content).map_err(|_| HostError::new("RECORD_SCHEMA_INVALID"))?; + let time = |value: i64| (0..=253402300799999).contains(&value); + if record.schema != 1 || record.kind != "task" || path(&record.id)? != path_value { + return Err(HostError::new("RECORD_SCHEMA_UNSUPPORTED")); + } + let data = &record.data; + if data.title.trim().is_empty() + || data.title.len() > 4096 + || data.description.len() > 262144 + || !matches!( + data.status.as_str(), + "todo" | "in_progress" | "done" | "cancelled" + ) + || !time(data.created_at_ms) + || !time(data.updated_at_ms) + || data.due_at_ms.is_some_and(|v| !time(v)) + || data.note_id.as_ref().is_some_and(|v| { + v.is_empty() + || v.len() > 128 + || !v + .bytes() + .all(|b| b.is_ascii_alphanumeric() || b"_-".contains(&b)) + }) + { + return Err(HostError::new("RECORD_DATA_INVALID")); + } + Ok(record) +} +impl Workspace { + pub(crate) fn normalize_record_links(&mut self) -> Result<()> { + for path in self.sync_paths()?.into_iter().filter(|v| allowed(v)) { + let bytes = std::fs::read(self.resolve(&path)?)?; + let original = validate(&path, &bytes)?; + if let Some(note_id) = original.data.note_id.as_ref() { + if let Ok(note_path) = self.path_for_id(note_id) { + if let Some(entry) = self.entry(¬e_path)? { + if &entry.file_id != note_id { + let mut record = original; + record.data.note_id = Some(entry.file_id); + self.write( + &path, + &crate::workspace::hash(&bytes), + &serde_json::to_vec(&record) + .map_err(|_| HostError::new("RECORD_SCHEMA_INVALID"))?, + "local", + )?; + } + } + } + } + } + Ok(()) + } + pub fn record_get(&mut self, id: &str) -> Result> { + let path = path(id)?; + if !self.resolve(&path)?.is_file() { + return Ok(None); + } + if std::fs::metadata(self.resolve(&path)?)?.len() > 1024 * 1024 { + return Err(HostError::new("RECORD_TOO_LARGE")); + } + let document = self.read(&path)?; + let mut record = validate(&path, document.content.as_bytes())?; + if let Some(note_id) = record.data.note_id.as_ref() { + if let Ok(path) = self.path_for_id(note_id) { + if let Some(entry) = self.entry(&path)? { + record.data.note_id = Some(entry.file_id); + } + } + } + Ok(Some( + json!({"record":record,"hash":document.entry.hash,"file_id":document.entry.file_id}), + )) + } + pub fn record_list(&mut self, offset: usize, limit: usize) -> Result { + if limit == 0 || limit > 1000 { + return Err(HostError::new("RECORD_LIMIT_INVALID")); + } + let paths = self + .sync_paths()? + .into_iter() + .filter(|path| allowed(path)) + .collect::>(); + let total = paths.len(); + let mut items = Vec::new(); + let mut bytes = 0; + for path in paths.into_iter().skip(offset).take(limit) { + let id = path + .strip_prefix("opennexus-records/v1/tasks/") + .and_then(|v| v.strip_suffix(".json")) + .ok_or_else(|| HostError::new("RECORD_ID_INVALID"))?; + if let Some(value) = self.record_get(id)? { + let size = serde_json::to_vec(&value) + .map_err(|_| HostError::new("RECORD_SCHEMA_INVALID"))? + .len(); + if bytes + size > 4 * 1024 * 1024 { + break; + } + bytes += size; + items.push(value); + } + } + Ok(json!({"items":items,"total":total})) + } + pub fn record_operation(&mut self, operation: &str) -> Result> { + self.recover()?; + let Some(receipt) = self.operation(operation)? else { + return Ok(None); + }; + let path = receipt["result"]["path"] + .as_str() + .ok_or_else(|| HostError::new("RECORD_OPERATION_PENDING"))?; + if !allowed(path) { + return Err(HostError::new("RECORD_OPERATION_DENIED")); + } + let bytes = self.payload(operation, &[])?; + let record = validate(path, &bytes)?; + Ok(Some( + json!({"record":record,"deleted":receipt["result"]["deleted"],"state":receipt["state"]}), + )) + } +} +#[cfg(test)] +mod tests { + use super::*; + fn fixture() -> Value { + json!({"schema":1,"kind":"task","id":"task_00000000000000000000000000000001","data":{"title":"Task","description":"","status":"todo","note_id":null,"due_at_ms":null,"created_at_ms":0,"updated_at_ms":0}}) + } + #[test] + fn whitelist_rejects_secret_unknown_fields_and_future_schema_before_journal() { + let root = tempfile::tempdir().unwrap(); + let mut ws = Workspace::open(root.path()).unwrap(); + let value = fixture(); + let path = path(value["id"].as_str().unwrap()).unwrap(); + for field in ["api_key", "token", "environment", "permissions", "extra"] { + let mut value = value.clone(); + value["data"][field] = json!("planted-secret"); + assert_eq!( + ws.write(&path, "", &serde_json::to_vec(&value).unwrap(), "local") + .unwrap_err() + .code, + "RECORD_SCHEMA_INVALID" + ); + } + let mut future = value.clone(); + future["schema"] = json!(2); + assert_eq!( + ws.write(&path, "", &serde_json::to_vec(&future).unwrap(), "remote") + .unwrap_err() + .code, + "RECORD_SCHEMA_UNSUPPORTED" + ); + assert!(!root.path().join(&path).exists()); + assert_eq!(ws.pending_count().unwrap(), 0); + let operation = uuid::Uuid::new_v4().to_string(); + ws.write_operation( + &path, + "", + &serde_json::to_vec(&value).unwrap(), + "local", + &operation, + ) + .unwrap(); + assert_eq!( + ws.record_operation(&operation).unwrap().unwrap()["record"], + value + ); + } +} diff --git a/frontend/src-tauri/src/sync_discovery.rs b/frontend/src-tauri/src/sync_discovery.rs index 68ab1f7..82ac11d 100644 --- a/frontend/src-tauri/src/sync_discovery.rs +++ b/frontend/src-tauri/src/sync_discovery.rs @@ -5,6 +5,9 @@ use std::{collections::HashSet, fs, io::Read, path::Path}; use uuid::Uuid; /// File transport only. Logical records receive their own versioned whitelist separately. pub fn allowed(path: &str) -> bool { + if crate::records::is_record(path) { + return crate::records::allowed(path); + } let parts: Vec<_> = path.split('/').collect(); if parts.iter().any(|part| { part.starts_with('.') @@ -100,6 +103,9 @@ impl Workspace { if bytes.len() > 104857600 { return Err(HostError::new("FILE_TOO_LARGE")); } + if crate::records::is_record(&path) { + crate::records::validate(&path, &bytes)?; + } let digest = hash(&bytes); let previous = self.entry(&path)?; let observed: Option<(String, String, bool)> = if let Some(entry) = &previous { diff --git a/frontend/src-tauri/src/sync_state.rs b/frontend/src-tauri/src/sync_state.rs index 8f2d169..f527d6d 100644 --- a/frontend/src-tauri/src/sync_state.rs +++ b/frontend/src-tauri/src/sync_state.rs @@ -140,6 +140,7 @@ impl Workspace { } pub fn sync_capture(&mut self, binding: &str) -> Result<()> { self.check_binding(binding)?; + self.normalize_record_links()?; 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)?)) @@ -165,6 +166,9 @@ impl Workspace { return Err(HostError::new("SYNC_SPOOL_CORRUPT")); } crate::payloads::verify(&self.sync_spool(&digest)?, &digest, size as u64)?; + if crate::records::is_record(&path) { + crate::records::validate(&path, &self.payload(&operation_id, &content)?)?; + } size } else { 0 diff --git a/frontend/src-tauri/src/workspace.rs b/frontend/src-tauri/src/workspace.rs index 498367f..01bbac0 100644 --- a/frontend/src-tauri/src/workspace.rs +++ b/frontend/src-tauri/src/workspace.rs @@ -258,10 +258,11 @@ impl Workspace { .map_err(|_| HostError::new("UNSAFE_PATH"))? .to_string_lossy() .replace('\\', "/"); - if relative - .split('/') - .any(|s| s.eq_ignore_ascii_case(".ainote") || s.eq_ignore_ascii_case(".git")) - || linked(&path)? + if relative.split('/').any(|s| { + s.eq_ignore_ascii_case(".ainote") + || s.eq_ignore_ascii_case(".git") + || s.eq_ignore_ascii_case("opennexus-records") + }) || linked(&path)? { continue; } @@ -407,6 +408,9 @@ impl Workspace { if Uuid::parse_str(operation_id).is_err() { return Err(HostError::new("OPERATION_ID_INVALID")); } + if crate::records::is_record(path) { + crate::records::validate(path, content)?; + } if content.len() > 100 * 1024 * 1024 { return Err(HostError::new("FILE_TOO_LARGE")); } @@ -526,6 +530,9 @@ impl Workspace { content: &[u8], origin: &str, ) -> Result<()> { + if crate::records::is_record(path) { + crate::records::validate(path, content)?; + } let target = self.resolve(path)?; let digest = hash(content); let current = if target.exists() { diff --git a/frontend/src-tauri/src/workspace_broker.rs b/frontend/src-tauri/src/workspace_broker.rs index 41ed305..9278b33 100644 --- a/frontend/src-tauri/src/workspace_broker.rs +++ b/frontend/src-tauri/src/workspace_broker.rs @@ -42,6 +42,28 @@ struct Mutation { kind: String, } +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct RecordRead { + vault_id: String, + id: String, +} +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct RecordDelete { + vault_id: String, + id: String, + expected: String, + operation_id: String, +} +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct RecordWrite { + vault_id: String, + record: Value, + expected: String, + operation_id: String, +} fn bound(ws: &Workspace, vault_id: &str) -> Result<(), String> { if ws.vault_id != vault_id { return Err("VAULT_PERMISSION_CHANGED".into()); @@ -54,6 +76,47 @@ fn decode(value: &Value) -> Result { pub fn dispatch(ws: &mut Workspace, request: &Value) -> Result { let params = &request["params"]; match request["rpc"].as_str().unwrap_or_default() { + "workspace.records.list" => { + let p: List = decode(params)?; + bound(ws, &p.vault_id)?; + ws.record_list(p.offset, p.limit).map_err(|e| e.code) + } + "workspace.records.get" => { + let p: RecordRead = decode(params)?; + bound(ws, &p.vault_id)?; + ws.record_get(&p.id) + .map(|v| v.unwrap_or(Value::Null)) + .map_err(|e| e.code) + } + "workspace.records.operation" => { + let p: Operation = decode(params)?; + bound(ws, &p.vault_id)?; + ws.record_operation(&p.operation_id) + .map(|v| v.unwrap_or(Value::Null)) + .map_err(|e| e.code) + } + "workspace.records.delete" => { + let p: RecordDelete = decode(params)?; + bound(ws, &p.vault_id)?; + let path = crate::records::path(&p.id).map_err(|e| e.code)?; + ws.mutate_operation("delete", &path, "", &p.expected, &p.operation_id) + .map_err(|e| e.code)?; + ws.record_operation(&p.operation_id) + .map(|v| v.unwrap_or(Value::Null)) + .map_err(|e| e.code) + } + "workspace.records.write" => { + let p: RecordWrite = decode(params)?; + bound(ws, &p.vault_id)?; + let path = + crate::records::path(p.record["id"].as_str().unwrap_or("")).map_err(|e| e.code)?; + let bytes = serde_json::to_vec(&p.record).map_err(|_| "RECORD_SCHEMA_INVALID")?; + ws.write_operation(&path, &p.expected, &bytes, "local", &p.operation_id) + .map_err(|e| e.code)?; + ws.record_operation(&p.operation_id) + .map(|v| v.unwrap_or(Value::Null)) + .map_err(|e| e.code) + } "workspace.list" => { let p: List = decode(params)?; bound(ws, &p.vault_id)?; diff --git a/frontend/src-tauri/tests/core_workspace.rs b/frontend/src-tauri/tests/core_workspace.rs index 3466a35..edbd84e 100644 --- a/frontend/src-tauri/tests/core_workspace.rs +++ b/frontend/src-tauri/tests/core_workspace.rs @@ -146,6 +146,122 @@ async fn real_core_notes_roundtrip_only_through_bound_host_and_confirm_commits() .await; assert_eq!(status, 409, "{denied}"); assert_eq!(denied["error"]["code"], "VAULT_PERMISSION_CHANGED"); + + let legacy_database = rusqlite::Connection::open( + temp.path() + .join("core/vault-state") + .join(&vault) + .join("core.sqlite3"), + ) + .unwrap(); + legacy_database.execute("INSERT INTO tasks VALUES ('task_00000000000000000000000000000002','Legacy task','','todo',?1,NULL,'2026-09-08T00:00:00+00:00','2026-09-08T00:00:00+00:00')",[file_id]).unwrap(); + let task_operation = uuid::Uuid::new_v4().to_string(); + let task_body = json!({"title":"Host task","note_id":file_id}); + let (status, task) = request( + &mut core, + "POST", + "/api/tasks", + &vault, + &task_operation, + Some(task_body.clone()), + ) + .await; + assert_eq!(status, 200, "{task}"); + let task_id = task["task_id"].as_str().unwrap(); + for _ in 0..20 { + let (status, replay) = request( + &mut core, + "POST", + "/api/tasks", + &vault, + &task_operation, + Some(task_body.clone()), + ) + .await; + assert_eq!(status, 200, "{replay}"); + assert_eq!(replay, task); + } + let (status, changed_request) = request( + &mut core, + "POST", + "/api/tasks", + &vault, + &task_operation, + Some(json!({"title":"changed","note_id":file_id})), + ) + .await; + assert_eq!(status, 409, "{changed_request}"); + let record_path = root.join(format!("opennexus-records/v1/tasks/{task_id}.json")); + assert!(record_path.is_file()); + let (status, task_updated) = request( + &mut core, + "PATCH", + &format!("/api/tasks/{task_id}"), + &vault, + &uuid::Uuid::new_v4().to_string(), + Some(json!({"status":"done"})), + ) + .await; + assert_eq!(status, 200, "{task_updated}"); + assert_eq!(task_updated["status"], "done"); + let (status, tasks) = request( + &mut core, + "GET", + "/api/tasks", + &vault, + &uuid::Uuid::new_v4().to_string(), + None, + ) + .await; + assert_eq!(status, 200, "{tasks}"); + assert_eq!(tasks["items"].as_array().unwrap().len(), 2); + let (status, task_denied) = request( + &mut core, + "GET", + "/api/tasks", + &uuid::Uuid::new_v4().to_string(), + &uuid::Uuid::new_v4().to_string(), + None, + ) + .await; + assert_eq!(status, 409, "{task_denied}"); + let deletion = uuid::Uuid::new_v4().to_string(); + for _ in 0..2 { + let (status, result) = request( + &mut core, + "DELETE", + &format!("/api/tasks/{task_id}"), + &vault, + &deletion, + None, + ) + .await; + assert_eq!(status, 200, "{result}"); + } + assert!(!record_path.exists()); + let source: String = legacy_database + .query_row( + "SELECT title FROM tasks WHERE task_id='task_00000000000000000000000000000002'", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(source, "Legacy task"); + legacy_database + .execute("DELETE FROM index_meta WHERE key='tasks_host_owned_v1'", []) + .unwrap(); + let (status, after_migration) = request( + &mut core, + "GET", + "/api/tasks", + &vault, + &uuid::Uuid::new_v4().to_string(), + None, + ) + .await; + assert_eq!(status, 200, "{after_migration}"); + assert_eq!(after_migration["items"].as_array().unwrap().len(), 1); + let (status, deleted) = request( &mut core, "DELETE", @@ -168,7 +284,7 @@ async fn real_core_notes_roundtrip_only_through_bound_host_and_confirm_commits() .await; assert_eq!(status, 200, "{search}"); assert!(!search.to_string().contains(file_id)); - assert_eq!(workspace.lock().unwrap().pending_count().unwrap(), 4); + assert_eq!(workspace.lock().unwrap().pending_count().unwrap(), 8); assert!(!temp .path() .join("core/unbound-vault/Core fixture.md") diff --git a/frontend/src-tauri/tests/sync_push.rs b/frontend/src-tauri/tests/sync_push.rs index ef4a894..fe1ef67 100644 --- a/frontend/src-tauri/tests/sync_push.rs +++ b/frontend/src-tauri/tests/sync_push.rs @@ -423,6 +423,67 @@ async fn actual_service_accepts_ordered_push_and_repeat_commit_without_duplicate .content, "only-local" ); + + let task_id = "task_00000000000000000000000000000001"; + let task_path = format!("opennexus-records/v1/tasks/{task_id}.json"); + let task_record = json!({"schema":1,"kind":"task","id":task_id,"data":{"title":"Synchronized task","description":"","status":"todo","note_id":first.file_id,"due_at_ms":null,"created_at_ms":0,"updated_at_ms":0}}); + workspace + .lock() + .unwrap() + .write( + &task_path, + "", + &serde_json::to_vec(&task_record).unwrap(), + "local", + ) + .unwrap(); + client.push_one(&workspace, &binding).await.unwrap(); + client_b.pull_page(&workspace_b, &binding_b).await.unwrap(); + { + let mut ws = workspace_b.lock().unwrap(); + let task = ws.record_get(task_id).unwrap().unwrap(); + assert_eq!(task["record"], task_record); + let mut next = task["record"].clone(); + next["data"]["status"] = json!("done"); + ws.write( + &task_path, + task["hash"].as_str().unwrap(), + &serde_json::to_vec(&next).unwrap(), + "local", + ) + .unwrap(); + } + client_b.push_one(&workspace_b, &binding_b).await.unwrap(); + client.pull_page(&workspace, &binding).await.unwrap(); + assert_eq!( + workspace + .lock() + .unwrap() + .record_get(task_id) + .unwrap() + .unwrap()["record"]["data"]["status"], + "done" + ); + { + let mut ws = workspace_b.lock().unwrap(); + let task = ws.record_get(task_id).unwrap().unwrap(); + ws.mutate_operation( + "delete", + &task_path, + "", + task["hash"].as_str().unwrap(), + &uuid::Uuid::new_v4().to_string(), + ) + .unwrap(); + } + client_b.push_one(&workspace_b, &binding_b).await.unwrap(); + client.pull_page(&workspace, &binding).await.unwrap(); + assert!(workspace + .lock() + .unwrap() + .record_get(task_id) + .unwrap() + .is_none()); // Kill the actual client process after each durable 10 MiB server offset, // before its response reaches the client. The next process must query offset. use sha2::{Digest, Sha256};