diff --git a/CHANGELOG.md b/CHANGELOG.md index 48047a872..0ded7f319 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -12,6 +12,13 @@ +## 未发布 / Unreleased + +- **sync**:保护 Native 离线编辑、删除冲突和迟到 ACK;重试保留 mutationId,避免旧更新复活已删除笔记。Protect durable native edits, deletion conflicts and late ACKs without changing mutation identity. +- **android**:复用 WebSocket 唤醒 Sync V2 Pull,提交后刷新界面,并通过前台周期 Pull 补偿丢失通知。Wake native Pull on notices/reconnect and refresh UI after committed Apply. +- **sync**:修复同页创建后删除、工作区 Snapshot 清理和冲突后继续编辑的数据保护问题。Protect pending descendants, attachment bytes and newer local edits during recovery. +- 验收范围与未完成平台见 `docs/sync-v2-reliability-validation.md`;这些改动尚未正式发布。See the validation record for evidence and remaining acceptance gaps; these changes are not released. + ## v1.5.1 - 2026-10-08 ### ✨ 新增 diff --git a/backend/src/routes/sync-v2.ts b/backend/src/routes/sync-v2.ts index 681f8b1f6..7d5addb2e 100644 --- a/backend/src/routes/sync-v2.ts +++ b/backend/src/routes/sync-v2.ts @@ -804,6 +804,8 @@ app.post("/push", async (c) => { if (parsed.entityType === "note" && typeof serverPayload?.version === "number") { serverVersion=serverPayload.version; serverPayload = withEncryptedBlocksSupport(serverPayload); + } else if (parsed.entityType === "note" && !serverPayload) { + serverPayload = { id: parsed.entityId, __delete: true }; } } } diff --git a/backend/src/sync/apply.ts b/backend/src/sync/apply.ts index 72afd6756..fb560ce94 100644 --- a/backend/src/sync/apply.ts +++ b/backend/src/sync/apply.ts @@ -248,6 +248,10 @@ function applyNote(db: Database.Database, input: ApplyMutationInput): number | n return nextVersion; } + // An update based on an existing revision is not a new entity. An old device + // must not resurrect a permanently deleted note, even after feed retention. + if (input.baseVersion !== undefined) throw new SyncError("VERSION_CONFLICT", "远端笔记已删除,请保留冲突或另存为新笔记"); + // 新建:客户端生成 UUID,离线也能创建,不依赖服务端分配 ID。 const version = Math.max(1, num(p.version, 1)); db.prepare(` diff --git a/backend/src/sync/engine.ts b/backend/src/sync/engine.ts index e51022ef1..bc90d1939 100644 --- a/backend/src/sync/engine.ts +++ b/backend/src/sync/engine.ts @@ -718,7 +718,9 @@ export class SyncEngine { const deletions: EngineRemotePayload[] = []; const wanted = new Set(); - for (const item of items) { + // Only the last operation per identity describes its state in Snapshot. + const latest = new Map(items.map((item) => [`${item.entityType}\u0000${item.entityId}`, item])); + for (const item of latest.values()) { if (item.operation === "delete") { deletions.push(item.entityType === "knowledge_tree_node" ? { @@ -779,8 +781,7 @@ export class SyncEngine { ); } - // 先应用 upsert 再应用 delete: - // 同一轮里若既有创建又有删除,删除应当是最终状态。 + // Apply only the final operation for each entity in this feed page. return [...upserts, ...deletions]; } @@ -990,7 +991,9 @@ export class SyncEngine { ${ids.length ? `AND id NOT IN (${placeholders})` : ""} AND NOT EXISTS (SELECT 1 FROM sync_outbox o WHERE o.profileId=? AND o.scopeKey=? AND o.entityType=? AND o.entityId=${table}.id AND o.status IN ('pending','inflight','failed')) - ${extra}`).run(scope.workspaceId,...ids,this.profileId,scope.scopeKey,entityType); + AND NOT EXISTS (SELECT 1 FROM sync_conflicts c WHERE c.profileId=? AND c.scopeKey=? + AND c.entityType=? AND c.entityId=${table}.id AND c.status='unresolved') + ${extra}`).run(scope.workspaceId,...ids,this.profileId,scope.scopeKey,entityType,this.profileId,scope.scopeKey,entityType); }; const removeComposite=(table:string,entityType:string,idSql:string,scopeSql:string)=>{ const ids=[...(seen.get(entityType) || [])]; @@ -998,19 +1001,24 @@ export class SyncEngine { this.db.prepare(`DELETE FROM ${table} WHERE ${scopeSql} ${ids.length ? `AND (${idSql}) NOT IN (${placeholders})` : ""} AND NOT EXISTS (SELECT 1 FROM sync_outbox o WHERE o.profileId=? AND o.scopeKey=? - AND o.entityType=? AND o.entityId=(${idSql}) AND o.status IN ('pending','inflight','failed'))`) - .run(scope.workspaceId,...ids,this.profileId,scope.scopeKey,entityType); + AND o.entityType=? AND o.entityId=(${idSql}) AND o.status IN ('pending','inflight','failed')) + AND NOT EXISTS (SELECT 1 FROM sync_conflicts c WHERE c.profileId=? AND c.scopeKey=? + AND c.entityType=? AND c.entityId=(${idSql}) AND c.status='unresolved')`) + .run(scope.workspaceId,...ids,this.profileId,scope.scopeKey,entityType,this.profileId,scope.scopeKey,entityType); }; runWithOutboxSuppressed(()=>runChangeFeedSuppressed(this.db,()=>this.db.transaction(()=>{ removeMissing("attachments","attachment"); removeComposite("favorites","favorite","favorites.userId || ':' || favorites.noteId","favorites.workspaceId=?"); removeComposite("note_tags","note_tag","note_tags.noteId || ':' || note_tags.tagId","note_tags.noteId IN (SELECT id FROM notes WHERE workspaceId=?)"); removeComposite("task_reminders","task_reminder","task_reminders.id","task_reminders.taskId IN (SELECT id FROM tasks WHERE workspaceId=?)"); - removeMissing("tasks","task"); + removeMissing("tasks","task","AND NOT EXISTS (SELECT 1 FROM task_reminders r WHERE r.taskId=tasks.id)"); removeMissing("diaries","diary"); removeMissing("mindmaps","mindmap"); - removeMissing("notes","note"); - removeMissing("notebooks","notebook","AND NOT EXISTS (SELECT 1 FROM notes n WHERE n.notebookId=notebooks.id)"); + removeMissing("notes","note",`AND NOT EXISTS (SELECT 1 FROM attachments a WHERE a.noteId=notes.id) + AND NOT EXISTS (SELECT 1 FROM note_tags nt WHERE nt.noteId=notes.id) + AND NOT EXISTS (SELECT 1 FROM favorites f WHERE f.noteId=notes.id)`); + removeMissing("notebooks","notebook",`AND NOT EXISTS (SELECT 1 FROM notes n WHERE n.notebookId=notebooks.id) + AND NOT EXISTS (SELECT 1 FROM notebooks child WHERE child.parentId=notebooks.id)`); removeMissing("tags","tag","AND NOT EXISTS (SELECT 1 FROM note_tags nt WHERE nt.tagId=tags.id)"); })())); } diff --git a/backend/src/sync/push.ts b/backend/src/sync/push.ts index 61192d920..a0b5f1e71 100644 --- a/backend/src/sync/push.ts +++ b/backend/src/sync/push.ts @@ -72,9 +72,8 @@ export function coalesceMutations(rows: SyncOutboxRow[]): CoalescedMutation[] { existing.operation = row.operation; existing.payload = parsePayload(row.payload); // baseVersion 保持最早那条:它才代表"这串修改的共同祖先"。 - if (existing.baseVersion === undefined && row.baseVersion !== null) { - existing.baseVersion = row.baseVersion; - } + // An absent first base means creation. Later local revisions must not turn + // that creation into an update of a server entity that does not exist yet. } return order.map((key) => byEntity.get(key) as CoalescedMutation); diff --git a/backend/tests/sync-v2-android-server.ts b/backend/tests/sync-v2-android-server.ts new file mode 100644 index 000000000..5e9dd8abd --- /dev/null +++ b/backend/tests/sync-v2-android-server.ts @@ -0,0 +1,49 @@ +// Explicit test helper, not a *.test.ts suite or production route. Loopback only. +import fs from "node:fs"; +import path from "node:path"; +import crypto from "node:crypto"; +import { Hono } from "hono"; +import { cors } from "hono/cors"; +import { serve } from "@hono/node-server"; + +async function main() { + const directory = process.env.SYNC_ACCEPTANCE_DIRECTORY; + if (!directory) throw new Error("SYNC_ACCEPTANCE_DIRECTORY required"); + fs.mkdirSync(directory, { recursive: true }); + process.env.DB_PATH = path.join(directory, "server.db"); + process.env.ELECTRON_USER_DATA = directory; + process.env.JWT_SECRET = "isolated-sync-acceptance-secret-2026"; + const { getDb } = await import("../src/db/schema"); + const { signLoginToken, verifyLoginToken } = await import("../src/lib/auth-security"); + const { attachRealtimeServer, broadcastToUser } = await import("../src/services/realtime"); + const { setSyncBroadcaster } = await import("../src/sync/notify"); + const { default: routes } = await import("../src/routes/sync-v2"); + const userId = "isolated-android-sync-user", db = getDb(); + db.prepare("INSERT OR IGNORE INTO users (id,username,passwordHash,createdAt,updatedAt) VALUES (?,?,'test-only','now','now')").run(userId, userId); + const fixtureToken = () => signLoginToken({ userId, username: userId, tokenVersion: 0 }); + const app = new Hono(); + app.use("*", cors({ origin: "https://localhost" })); + app.get("/acceptance/config", (c) => c.json({ userId, token: fixtureToken(), serverUrl: "http://127.0.0.1:47831" })); + app.use("/api/*", async (c, next) => { + const identity = verifyLoginToken(c.req.header("Authorization")?.replace(/^Bearer /, "") || ""); + if (!identity) return c.json({ error: "UNAUTHORIZED" }, 401); + c.req.raw.headers.set("X-User-Id", identity.userId); + await next(); + }); + app.route("/api/sync/v2", routes); + app.get("/acceptance/note/:id", (c) => c.json(db.prepare("SELECT * FROM notes WHERE id=? AND userId=?").get(c.req.param("id"), userId) || null)); + app.post("/acceptance/edit/:id", async (c) => { + const note = db.prepare("SELECT * FROM notes WHERE id=? AND userId=?").get(c.req.param("id"), userId) as Record | undefined; + if (!note) return c.json({ error: "NOT_FOUND" }, 404); + const { text } = await c.req.json<{text:string}>(); + const content = note.contentFormat === "markdown" ? text : JSON.stringify({ type: "doc", content: [{ type: "paragraph", content: [{ type: "text", text }] }] }); + return app.request("http://localhost/api/sync/v2/push", { method: "POST", headers: { Authorization: `Bearer ${fixtureToken()}`, "Content-Type": "application/json" }, + body: JSON.stringify({ scopeKey: "personal", deviceId: "http-client", mutations: [{ mutationId: crypto.randomUUID(), entityType: "note", entityId: note.id, operation: "upsert", baseVersion: note.version, + payload: { ...note, content, contentText: text } }] }) }); + }); + const server = serve({ fetch: app.fetch, hostname: "127.0.0.1", port: 47831 }); + attachRealtimeServer(server); + setSyncBroadcaster((recipient, notice) => broadcastToUser(recipient, notice as never)); + fs.writeFileSync(path.join(directory, "ready"), "ready"); +} +void main(); diff --git a/backend/tests/sync-v2-engine.test.ts b/backend/tests/sync-v2-engine.test.ts index 61e0ce0e4..ba7396f1f 100644 --- a/backend/tests/sync-v2-engine.test.ts +++ b/backend/tests/sync-v2-engine.test.ts @@ -706,6 +706,24 @@ test("无变更时也推进游标并 ACK,避免重复扫描", async () => { assert.deepEqual(remote.ackCalls, [42]); }); +test("同一批创建后删除只应用最终删除,不等待已不存在的 Snapshot", async () => { + resetSyncTables(); + const db = getDb(); + const { engine, remote, profileId } = createEngine(); + remote.serverSequence = 43; + remote.changesQueue.push({ + serverSequence: 43, nextSequence: 43, hasMore: false, resetRequired: false, + items: [ + { sequence: 42, entityType: "note", entityId: "created-then-deleted", operation: "upsert" }, + { sequence: 43, entityType: "note", entityId: "created-then-deleted", operation: "delete" }, + ], + }); + const status = await engine.syncOnce(); + assert.notEqual(status.state, "error"); + assert.equal(sync.getSyncState(db, profileId)?.lastSequence, 43); + assert.deepEqual(remote.ackCalls, [43]); +}); + test("Change Feed upsert 缺少 Snapshot payload 时不推进游标或 ACK", async () => { resetSyncTables(); const db = getDb(); diff --git a/backend/tests/sync-v2-mobile-multidevice.test.ts b/backend/tests/sync-v2-mobile-multidevice.test.ts new file mode 100644 index 000000000..c57f94e77 --- /dev/null +++ b/backend/tests/sync-v2-mobile-multidevice.test.ts @@ -0,0 +1,379 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { randomUUID } from "node:crypto"; +import Database from "better-sqlite3"; +import { Hono } from "hono"; +import { serve } from "@hono/node-server"; +import { NativeLocalRepository } from "../../frontend/src/lib/nativeLocalRepository"; +import { MobileSyncEngine } from "../../frontend/src/lib/mobileSyncEngine"; +import type { NativeDatabase } from "../../frontend/src/lib/nativeDatabase"; +import { nativeReceiptKey, readNativeNoteSyncReceipt } from "../../frontend/src/lib/mobileNoteSyncReceipt"; + +const tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), "nowen-sync-multidevice-")); +process.env.DB_PATH = path.join(tmpDir, "server.db"); +process.env.ELECTRON_USER_DATA = tmpDir; +process.env.NOWEN_LOCAL_FIRST_SYNC_V2 = "1"; +const browser = new EventTarget(); +Object.assign(globalThis, { window: browser }); +const scope = { scopeKey: "personal", workspaceId: null, workspaceName: null, role: "owner", canWrite: true, accessFingerprint: "personal:owner" }; +const USER = "multidevice-user"; +let serverDb: Database.Database; +let closeServerDb: () => void; +let httpServer: ReturnType; +let serverUrl: string; +const clients: Client[] = []; +let afterPush: ((response: Response) => Promise) | undefined; + +// The real native schema, with independent on-disk SQLite databases. The adapter +// serializes transactions just like NativeDatabaseImpl; no shared client state. +function openClientDb(file: string): NativeDatabase { + const raw = new Database(file); + raw.pragma("foreign_keys = ON"); + const source = fs.readFileSync(path.resolve(__dirname, "../../frontend/src/lib/nativeDatabase.ts"), "utf8"); + const scopeCheck = source.match(/const SCOPE_CHECK = `([\s\S]*?)`;/)![1]; + const entityTypes = source.match(/const ENTITY_TYPE_CHECK = `([\s\S]*?)`;/)![1]; + for (const [, sql] of source.matchAll(/`(CREATE TABLE IF NOT EXISTS [\s\S]*?)`/g)) { + raw.exec(sql.replaceAll("${SCOPE_CHECK}", scopeCheck).replaceAll("${ENTITY_TYPE_CHECK}", entityTypes)); + } + let queue = Promise.resolve(); + const schedule = (work: () => Promise): Promise => { + const next = queue.then(work, work); + queue = next.then(() => undefined, () => undefined); + return next; + }; + const tx: NativeDatabase = { + async run(sql, values = []) { const result = raw.prepare(sql).run(...values); return { changes: result.changes }; }, + async query(sql: string, values: unknown[] = []) { return raw.prepare(sql).all(...values) as T[]; }, + async transaction(work) { return work(tx); }, + async close() { raw.close(); }, + }; + return { + run: (sql, values) => schedule(() => tx.run(sql, values)), + query: (sql: string, values?: unknown[]) => schedule(() => tx.query(sql, values)), + transaction: (work) => schedule(async () => { + raw.exec("BEGIN IMMEDIATE"); + try { const result = await work(tx); raw.exec("COMMIT"); return result; } + catch (error) { raw.exec("ROLLBACK"); throw error; } + }), + close: () => schedule(() => tx.close()), + }; +} + +interface Client { + db: NativeDatabase; + file: string; + repository: NativeLocalRepository; + engine: MobileSyncEngine; + profileId: string; + removedFiles: string[]; + sync(): Promise; +} +async function client(name: string): Promise { + const file = path.join(tmpDir, `${name}.db`); + const db = openClientDb(file); + const profileId = `profile-${name}`; + await db.run(`INSERT OR IGNORE INTO sync_profiles (id,name,serverUrl,remoteUserId,enabled,createdAt,updatedAt) + VALUES (?,?,?, ?,1,'now','now')`, [profileId, name, serverUrl, USER]); + await db.run("INSERT OR IGNORE INTO sync_devices (profileId,deviceId,platform,createdAt) VALUES (?,?,'android','now')", [profileId, name]); + const repository = new NativeLocalRepository({ db, attachments: {} as never, accountId: name, userId: USER, getScopeKey: () => "personal" }); + const removedFiles: string[] = []; + const engine = new MobileSyncEngine({ db, attachments: { remove: async (id: string) => { removedFiles.push(id); } } as never, profileId, deviceId: name, serverUrl, token: USER, userId: USER }); + const result = { db, file, repository, engine, profileId, removedFiles, async sync() { + engine.start(); + try { await engine.syncOnce(); } finally { engine.stop(); } + } }; + clients.push(result); + return result; +} +async function pending(c: Client) { return c.db.query<{ mutationId: string; baseVersion: number; payload: string }>("SELECT mutationId,baseVersion,payload FROM sync_outbox WHERE profileId=? ORDER BY createdAt,rowid", [c.profileId]); } +async function conflicts(c: Client) { return c.db.query<{ localPayload: string; remotePayload: string }>("SELECT localPayload,remotePayload FROM sync_conflicts WHERE profileId=? AND status='unresolved'", [c.profileId]); } + +test.before(async () => { + const schema = await import("../src/db/schema"); + const routes = await import("../src/routes/sync-v2"); + serverDb = schema.getDb(); closeServerDb = schema.closeDb; + for (const user of [USER, "other-user"]) serverDb.prepare("INSERT INTO users (id,username,passwordHash,createdAt,updatedAt) VALUES (?,?,'x','now','now')").run(user, user); + const app = new Hono(); + // Route tests use fixture identities; production auth remains in index.ts. + app.use("*", async (c, next) => { + c.req.raw.headers.set("X-User-Id", c.req.header("Authorization")?.replace("Bearer ", "") || ""); + await next(); + if (c.req.path.endsWith("/push") && afterPush) c.res = await afterPush(c.res); + }); + app.route("/api/sync/v2", routes.default); + httpServer = serve({ fetch: app.fetch, hostname: "127.0.0.1", port: 0 }); + await new Promise((resolve) => httpServer.once("listening", resolve)); + const address = httpServer.address(); assert.ok(address && typeof address !== "string"); + serverUrl = `http://127.0.0.1:${address.port}`; +}); +test.after(async () => { + for (const c of clients) { c.engine.stop(); await c.db.close().catch(() => undefined); } + await new Promise((resolve) => httpServer.close(() => resolve())); + closeServerDb(); + fs.rmSync(tmpDir, { recursive: true, force: true }); +}); + +test("workspace Snapshot pruning preserves pending descendants, unresolved conflicts and local attachment bytes", async () => { + const c = await client("prune-native"); + const key = "workspace:prune", book = randomUUID(), pendingNote = randomUUID(), conflictedNote = randomUUID(), attachment = randomUUID(); + await c.db.run("INSERT INTO notebooks (id,scopeKey,workspaceId,userId,name,createdAt,updatedAt) VALUES (?,?,'prune',?,'Book','now','now')", [book,key,USER]); + for (const id of [pendingNote, conflictedNote]) await c.db.run("INSERT INTO notes (id,scopeKey,workspaceId,userId,notebookId,title,content,createdAt,updatedAt) VALUES (?,?,'prune',?,?,'Title','retained','now','now')", [id,key,USER,book]); + await c.db.run("INSERT INTO sync_outbox (id,mutationId,profileId,deviceId,scopeKey,entityType,entityId,operation,payload,createdAt) VALUES (?,?,?,?,?,'note',?,'upsert','{}','now')", [randomUUID(),randomUUID(),c.profileId,"prune-native",key,pendingNote]); + await c.db.run("INSERT INTO sync_conflicts (id,profileId,scopeKey,entityType,entityId,localPayload,remotePayload,createdAt) VALUES (?,?,?,'note',?,'{}','{}','now')", [randomUUID(),c.profileId,key,conflictedNote]); + await c.db.run("INSERT INTO attachments (id,scopeKey,workspaceId,userId,noteId,filename,mimeType,size,localPath,available,transferStatus,createdAt,updatedAt) VALUES (?,?,'prune',?,?,'a.png','image/png',4,'fixture',1,'pending_upload','now','now')", [attachment,key,USER,pendingNote]); + await (c.engine as unknown as { pruneWorkspaceSnapshot(s: unknown, seen: Map>): Promise }).pruneWorkspaceSnapshot({ ...scope, scopeKey: key, workspaceId: "prune" }, new Map()); + assert.equal((await c.db.query("SELECT id FROM notes WHERE scopeKey=?", [key])).length, 2); + assert.equal((await c.db.query("SELECT id FROM notebooks WHERE scopeKey=?", [key])).length, 1); + assert.equal((await c.db.query("SELECT id FROM attachments WHERE scopeKey=?", [key])).length, 1); + assert.deepEqual(c.removedFiles, []); +}); + +test("create then permanent delete between Pulls advances the cursor without a missing Snapshot error", async () => { + const a = await client("fold-a"), b = await client("fold-b"); + const book = randomUUID(), id = randomUUID(); + await a.repository.notebooks.create({ id: book, name: "Book" }); + await a.sync(); await b.sync(); + await a.repository.notes.create({ id, notebookId: book, content: "temporary" }); + await a.sync(); + await a.repository.notes.remove(id); + await a.sync(); + await b.sync(); + assert.equal((await b.repository.notes.get(id)), null); + const state = (await b.db.query<{lastError:string|null;lastSequence:number}>("SELECT lastError,lastSequence FROM sync_state WHERE profileId=?", [b.profileId]))[0]; + assert.equal(state.lastError, null); + assert.equal(state.lastSequence, serverDb.prepare("SELECT MAX(sequence) AS seq FROM sync_changes_v2 WHERE userId=?").get(USER).seq); +}); + +test("three independent SQLite clients converge bidirectionally after consecutive Markdown and rich-text edits", async () => { + for (const format of ["markdown", "tiptap-json"]) { + const a = await client(`a-${format}`), b = await client(`b-${format}`), c = await client(`c-${format}`); + const book = randomUUID(), id = randomUUID(); + await a.repository.notebooks.create({ id: book, name: "Book" }); + const content = format === "markdown" ? "first" : JSON.stringify({ type: "doc", content: [{ type: "paragraph", content: [{ type: "text", text: "first" }] }] }); + await a.repository.notes.create({ id, notebookId: book, content, contentFormat: format }); + await a.sync(); await b.sync(); await c.sync(); + assert.equal((await b.repository.notes.get(id))?.content, content); + const base = (await a.repository.notes.get(id))!.version; + await a.repository.notes.update(id, { title: "edit 1" }); + await a.repository.notes.update(id, { title: "edit 2" }); + await a.sync(); + assert.equal((await conflicts(a)).length, 0, "sequential local writes are descendants, not concurrent conflicts"); + assert.equal((await pending(a)).length, 0); + await b.sync(); await c.sync(); + assert.equal((await b.repository.notes.get(id))?.title, "edit 2"); + assert.ok((await a.repository.notes.get(id))!.version > base); + await b.repository.notes.update(id, { title: "from B" }); + await b.sync(); await a.sync(); await c.sync(); + assert.equal((await a.repository.notes.get(id))?.title, "from B"); + assert.equal((await c.repository.notes.get(id))?.version, (await a.repository.notes.get(id))?.version); + } +}); + +test("remote permanent deletion preserves offline body, pending mutation and recoverable conflict", async () => { + const a = await client("delete-a"), b = await client("delete-b"); + const book = randomUUID(), id = randomUUID(); + await a.repository.notebooks.create({ id: book, name: "Book" }); + await a.repository.notes.create({ id, notebookId: book, content: "base", contentFormat: "markdown" }); + await a.sync(); await b.sync(); + await b.repository.notes.update(id, { content: "offline recovery body" }); + await a.repository.notes.remove(id); await a.sync(); + const pull = b.engine as unknown as { pull(scope: typeof scope): Promise }; + await pull.pull(scope); + assert.equal((await b.repository.notes.get(id))?.content, "offline recovery body"); + assert.equal((await pending(b)).length, 1); + const saved = await conflicts(b); + assert.equal(saved.length, 1); + assert.equal(JSON.parse(saved[0].localPayload).content, "offline recovery body"); + assert.equal(JSON.parse(saved[0].remotePayload).__delete, true); +}); + +test("offline create survives database close/reopen and immediate local read before uploading", async () => { + const a = await client("restart-a"); + const book = randomUUID(), id = randomUUID(); + await a.repository.notebooks.create({ id: book, name: "Offline book" }); + await a.repository.notes.create({ id, notebookId: book, content: "durable", contentFormat: "markdown" }); + const mutations = (await pending(a)).map((row) => row.mutationId); + assert.equal((await a.repository.notes.get(id))?.content, "durable"); + await a.db.close(); + const restarted = await client("restart-a"); + assert.equal((await restarted.repository.notes.get(id))?.content, "durable"); + assert.deepEqual((await pending(restarted)).map((row) => row.mutationId), mutations); + await restarted.sync(); + const b = await client("restart-b"); await b.sync(); + assert.equal((await b.repository.notes.get(id))?.content, "durable"); +}); + +test("missing or invalid successful ACK version cannot claim cloud confirmation", async () => { + const a = await client("receipt-a"); + for (const version of [null, 0, -1, "4"]) { + await a.db.run("INSERT OR REPLACE INTO native_runtime_meta (key,value,updatedAt) VALUES (?,?,'now')", [nativeReceiptKey(a.profileId, "test-note"), JSON.stringify({ mutationId: "m", serverVersion: version })]); + assert.equal((await readNativeNoteSyncReceipt(a.db, a.profileId, "test-note")).phase, "unverified"); + } +}); + +test("an old offline update cannot resurrect a permanently deleted server note", async () => { + const a = await client("resurrect-a"), b = await client("resurrect-b"); + const book = randomUUID(), id = randomUUID(); + await a.repository.notebooks.create({ id: book, name: "Book" }); + await a.repository.notes.create({ id, notebookId: book, content: "base", contentFormat: "markdown" }); + await a.sync(); await b.sync(); + await b.repository.notes.update(id, { content: "offline retained" }); + await a.repository.notes.remove(id); await a.sync(); + await b.sync(); + assert.equal(serverDb.prepare("SELECT id FROM notes WHERE id=?").get(id), undefined); + assert.equal((await b.repository.notes.get(id))?.content, "offline retained"); + assert.equal((await conflicts(b)).length, 1); +}); + +test("apply failure rolls back the entire batch and leaves cursor and ACK unchanged", async () => { + const a = await client("apply-a"), b = await client("apply-b"); + const book = randomUUID(), id = randomUUID(); + await a.repository.notebooks.create({ id: book, name: "Book" }); + await a.repository.notes.create({ id, notebookId: book, content: "base", contentFormat: "markdown" }); + await a.sync(); await b.sync(); + const state = await b.db.query("SELECT * FROM sync_state"); + await a.repository.notes.update(id, { content: "after failure" }); await a.sync(); + const run = b.db.transaction; + b.db.transaction = (work) => run(async (tx) => work({ ...tx, run: async (sql, values) => { + if (sql.includes("INSERT INTO notes")) throw new Error("SQLITE_FULL fixture"); + return tx.run(sql, values); + } })); + await assert.rejects((b.engine as unknown as { pull(s: typeof scope): Promise }).pull(scope), /SQLITE_FULL/); + assert.equal((await b.repository.notes.get(id))?.content, "base"); + assert.deepEqual(await b.db.query("SELECT * FROM sync_state"), state); + b.db.transaction = run; + await b.sync(); + assert.equal((await b.repository.notes.get(id))?.content, "after failure"); +}); + +test("successful push replay after client ACK crash uses the same mutationId exactly once", async () => { + const a = await client("ack-crash-a"); + const book = randomUUID(), id = randomUUID(); + await a.repository.notebooks.create({ id: book, name: "Book" }); + await a.repository.notes.create({ id, notebookId: book, content: "durable ACK", contentFormat: "markdown" }); + const before = await pending(a); + const run = a.db.transaction; + a.db.transaction = (work) => run(async (tx) => work({ ...tx, run: async (sql, values) => { + if (sql.includes("DELETE FROM sync_outbox WHERE mutationId")) throw new Error("ACK crash fixture"); + return tx.run(sql, values); + } })); + await assert.rejects((a.engine as unknown as { push(s: typeof scope): Promise }).push(scope), /ACK crash/); + assert.equal(serverDb.prepare("SELECT content FROM notes WHERE id=?").get(id)?.content, "durable ACK"); + assert.deepEqual((await pending(a)).map((row) => row.mutationId), before.map((row) => row.mutationId)); + a.db.transaction = run; + await (a.engine as unknown as { push(s: typeof scope): Promise }).push(scope); + assert.equal((await pending(a)).length, 0); + for (const row of before) assert.equal((serverDb.prepare("SELECT COUNT(*) AS count FROM sync_v2_applied_mutations WHERE mutationId=?").get(row.mutationId) as {count:number}).count, 1); +}); + +test("overlapping native writes read the latest revision inside the write transaction", async () => { + const a = await client("overlap-a"); + const book = randomUUID(), id = randomUUID(); + await a.repository.notebooks.create({ id: book, name: "Book" }); + await a.repository.notes.create({ id, notebookId: book, content: "base", contentFormat: "markdown" }); + await a.sync(); + await Promise.all([ + a.repository.notes.update(id, { title: "new title" }), + a.repository.notes.update(id, { content: "new body" }), + ]); + const note = await a.repository.notes.get(id); + assert.equal(note?.title, "new title"); assert.equal(note?.content, "new body"); + const rows = await pending(a); + assert.equal(rows[1].baseVersion, JSON.parse(rows[0].payload).version); + await a.sync(); + assert.equal((await conflicts(a)).length, 0); +}); + +test("concurrent devices preserve both versions, and keep-local resolution establishes the remote base", async () => { + const a = await client("conflict-a"), b = await client("conflict-b"); + const book = randomUUID(), id = randomUUID(); + await a.repository.notebooks.create({ id: book, name: "Book" }); + await a.repository.notes.create({ id, notebookId: book, content: "base", contentFormat: "markdown" }); + await a.sync(); await b.sync(); + await b.repository.notes.update(id, { content: "B version" }); + for (let index = 0; index < 4; index++) { await a.repository.notes.update(id, { content: "A version" }); await a.sync(); } + await b.sync(); + const rows = await b.db.query<{id:string;localPayload:string;remotePayload:string}>("SELECT * FROM sync_conflicts WHERE profileId=? AND status='unresolved'", [b.profileId]); + assert.equal(rows.length, 1); + assert.equal(JSON.parse(rows[0].localPayload).content, "B version"); + assert.equal(JSON.parse(rows[0].remotePayload).content, "A version"); + await b.repository.notes.update(id, { content: "B edited after conflict" }); + await b.sync(); + const fork = await b.engine.forkLocalConflict(rows[0].id, "local"); + assert.equal((await b.repository.notes.get(fork))?.content, "B edited after conflict"); + await b.engine.resolveLocalConflict(rows[0].id, "keep-local"); + await b.sync(); + assert.equal((await b.repository.notes.get(id))?.version, (serverDb.prepare("SELECT version FROM notes WHERE id=?").get(id) as {version:number}).version); + await b.repository.notes.update(id, { title: "after resolution" }); await b.sync(); await a.sync(); + assert.equal((await conflicts(b)).length, 0); + assert.equal((await a.repository.notes.get(id))?.content, "B edited after conflict"); +}); + +test("two edits committed during a delayed ACK survive own echo and subsequent push", async () => { + const a = await client("delayed-a"), b = await client("delayed-b"); + const book = randomUUID(), id = randomUUID(); + await a.repository.notebooks.create({ id: book, name: "Book" }); + await a.repository.notes.create({ id, notebookId: book, content: "base", contentFormat: "markdown" }); + await a.sync(); await b.sync(); + await a.repository.notes.update(id, { content: "sent version" }); + let release!: () => void, received!: () => void; + const gate = new Promise((resolve) => { release = resolve; }); + const committed = new Promise((resolve) => { received = resolve; }); + afterPush = async (response) => { received(); await gate; return response; }; + const internal = a.engine as unknown as { push(s: typeof scope): Promise; pull(s: typeof scope): Promise }; + try { + const pushing = internal.push(scope); + await committed; + await a.repository.notes.update(id, { content: "newer local 1" }); + await a.repository.notes.update(id, { content: "newer local 2" }); + release(); await pushing; + afterPush = undefined; + await internal.pull(scope); + assert.equal((await conflicts(a)).length, 0, "multiple later edits must not turn own ACK into a conflict"); + assert.equal((await a.repository.notes.get(id))?.content, "newer local 2"); + assert.equal((await pending(a)).length, 2); + await a.sync(); await b.sync(); + assert.equal((await conflicts(a)).length, 0); + assert.equal((await b.repository.notes.get(id))?.content, "newer local 2"); + } finally { release(); afterPush = undefined; } +}); + +test("invalid note ACK leaves the durable mutation available for an idempotent retry", async () => { + const a = await client("invalid-ack-a"); + const book = randomUUID(), id = randomUUID(); + await a.repository.notebooks.create({ id: book, name: "Book" }); + await a.repository.notes.create({ id, notebookId: book, content: "retained", contentFormat: "markdown" }); + const before = await pending(a); + afterPush = async (response) => { + const json = await response.json(); + for (const result of json.results) delete result.version; + return Response.json(json); + }; + try { + await assert.rejects((a.engine as unknown as {push(s:typeof scope):Promise}).push(scope), /ACK/); + assert.deepEqual((await pending(a)).map((row) => row.mutationId), before.map((row) => row.mutationId)); + assert.equal((await readNativeNoteSyncReceipt(a.db, a.profileId, id)).phase, "pending"); + } finally { afterPush = undefined; } + await a.sync(); + assert.equal((await pending(a)).length, 0); + assert.equal((await a.repository.notes.get(id))?.content, "retained"); +}); + +test("a failed Outbox insertion rolls back note creation in the same SQLite transaction", async () => { + const a = await client("outbox-failure-a"); + const book = randomUUID(), id = randomUUID(); + await a.repository.notebooks.create({ id: book, name: "Book" }); + const transaction = a.db.transaction; + a.db.transaction = (work) => transaction(async (tx) => work({ ...tx, run: async (sql, values) => { + if (sql.includes("INSERT INTO sync_outbox")) throw new Error("Outbox full fixture"); + return tx.run(sql, values); + } })); + await assert.rejects(a.repository.notes.create({ id, notebookId: book, content: "must not falsely save", contentFormat: "markdown" }), /Outbox full/); + assert.equal(await a.repository.notes.get(id), null); + a.db.transaction = transaction; + await a.repository.notes.create({ id, notebookId: book, content: "saved after retry", contentFormat: "markdown" }); + assert.equal((await a.repository.notes.get(id))?.content, "saved after retry"); +}); diff --git a/docs/architecture/data-consistency-contract.matrix.json b/docs/architecture/data-consistency-contract.matrix.json index 3e0d94936..11e347a7b 100644 --- a/docs/architecture/data-consistency-contract.matrix.json +++ b/docs/architecture/data-consistency-contract.matrix.json @@ -254,6 +254,29 @@ "case": "updates the real editor badge after a committed Android edit" } ] + }, + { + "id": "DCC-012", + "title": "Native independent SQLite clients preserve durable edits across sync failures", + "status": "covered", + "checks": [ + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "three independent SQLite clients converge bidirectionally" }, + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "remote permanent deletion preserves offline body" }, + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "offline create survives database close/reopen" }, + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "missing or invalid successful ACK version cannot claim cloud confirmation" }, + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "an old offline update cannot resurrect" }, + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "apply failure rolls back the entire batch" }, + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "successful push replay after client ACK crash" }, + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "overlapping native writes read the latest revision" }, + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "concurrent devices preserve both versions" }, + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "two edits committed during a delayed ACK" }, + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "invalid note ACK leaves the durable mutation" }, + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "a failed Outbox insertion rolls back note creation" }, + { "file": "frontend/src/lib/__tests__/mobileSyncEndpoint.test.ts", "case": "rejects a late push ACK when the account changes" }, + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "workspace Snapshot pruning preserves pending descendants" }, + { "file": "backend/tests/sync-v2-mobile-multidevice.test.ts", "case": "create then permanent delete between Pulls" }, + { "file": "frontend/src/lib/encryptedNotes/__tests__/nativeDatabaseInitialization.test.ts", "case": "external UI reads wait for an in-progress Apply transaction" } + ] } ] } diff --git a/docs/architecture/data-consistency-contract.md b/docs/architecture/data-consistency-contract.md index 3b5a180ee..47649fdda 100644 --- a/docs/architecture/data-consistency-contract.md +++ b/docs/architecture/data-consistency-contract.md @@ -45,6 +45,8 @@ ## 3. 关键时序约束 +- **DCC-012 — Native 多库故障恢复必须保留实体身份和未确认数据。** 三个独立客户端 SQLite 与独立服务端通过真实 Sync V2 HTTP 路由验证双向收敛;远端永久删除不能覆盖离线正文,旧版本更新不得复活已删除笔记;本机并行保存必须在写事务内读取版本;ACK 缺失/非法时保留 mutationId,ACK 延迟期间的新修改继续保留;冲突双方和解决后的版本基线均可验证。这是自动化 SQLite/API 证据,不代表 Android 设备或生产网络已验收。 + ### A. 本地写入 `编辑请求 → 本地事务提交(含同步意图) → 允许显示“已保存到本机” → 后台推送 → 服务端确认对应 mutation / 版本 → 允许显示“云端已确认”`。 diff --git a/docs/sync-realtime.md b/docs/sync-realtime.md index 558b773ab..fbc823642 100644 --- a/docs/sync-realtime.md +++ b/docs/sync-realtime.md @@ -2,6 +2,14 @@ 本文档描述 Nowen Note 的 WebSocket 实时同步架构,涵盖连接管理、事件分发、删除/回收站同步流程、前端行为以及离线兜底策略。 +### Sync V2 / Android Native(release/v1.5.2 开发中) + +Android 登录态复用现有 `/ws`:`sync.changed`、连接恢复和前台恢复唤醒同一个 `MobileSyncEngine`,30 秒前台补偿 Pull 处理丢失通知。通知不携带正文,也不替代 Change Feed;Native SQLite Apply 提交后才广播 `nowen:mobile-sync-applied`,刷新列表、知识树和打开正文。Android 不再由 `OfflineSyncRuntime` 启动另一套 Web 离线队列。系统冻结/强杀时不承诺持续后台实时同步。 + +Desktop uses the embedded backend Sync V2 engine; Android uses its native SQLite engine. Web/mobile browsers still use REST plus IndexedDB/offline queue and are not a unified transactional Local-first implementation. iOS/HarmonyOS equivalence has not been accepted. A note ACK confirms the server revision, not receipt by every device or completion of every attachment transfer. + +平台与实体矩阵、根因和执行证据见 [审计](sync-v2-reliability-audit.md) 与 [验收记录](sync-v2-reliability-validation.md)。当前改动仍需通过剩余验收,不能据此宣称生产治理完成。 + --- ## 1. WebSocket 连接模型 diff --git a/docs/sync-v2-reliability-audit.md b/docs/sync-v2-reliability-audit.md new file mode 100644 index 000000000..998bef54a --- /dev/null +++ b/docs/sync-v2-reliability-audit.md @@ -0,0 +1,85 @@ +# Sync V2 reliability audit / 多端同步可靠性审计 + +审计日期:2026-10-11。基准:`release/v1.5.2`,`d95d3787`。本文区分源码能力、自动化验证和设备验收;目标延迟不是实测结果。 + +## Phase 0:真实架构和平台路径 + +```mermaid +flowchart LR + Desktop[Electron renderer] --> REST[Embedded Hono REST] + REST --> SQLite[Desktop SQLite + transactional Outbox triggers] + SQLite --> Engine[backend SyncEngine] + Android[Native UI bridge] --> Native[NativeLocalRepository + SQLite transaction] + Native --> Mobile[MobileSyncEngine] + Web[Browser / Harmony ArkWeb] --> Remote[Remote REST + offline queue/cache] + Engine --> Push[Sync V2 Push] + Mobile --> Push + Push --> Server[Server SQLite transaction + mutation dedup + change feed] + Remote --> Server + Server --> Notice[WebSocket sync.changed notification] + Notice --> Pull[Sync V2 changes + snapshot + ACK] + Pull --> Engine + Pull --> Mobile + Engine --> SQLite + Mobile --> Native +``` + +| 平台 | 实际业务保存 | 同步路径与边界 | +| --- | --- | --- | +| Windows/macOS/Linux Electron | 内嵌 Hono → SQLite;触发器捕获 Outbox | `backend/src/sync/engine.ts`,Push/Pull/冲突/附件;各操作系统包没有在本次审计中验收 | +| Web,包括手机浏览器 | 在线 REST 写服务器;失败后 `offlineQueue.ts` + `localStore.ts` | **Remote-first**,`syncEngine.ts`/`offlineRead.ts` 是缓存和离线回放,不是 Native Sync V2;缓存与队列尚不是单个 IndexedDB 业务事务 | +| Android | `mobileLocalFirstBridge.ts` → `nativeLocalRepository.ts` → `nativeDatabase.ts` | 笔记与 Outbox 在同个 SQLite 事务;`mobileSyncEngine.ts` 使用十类旧协议实体;登录态知识树为服务端优先,正文为本地优先 | +| iOS | `configureRuntime()` 使用 `Capacitor.isNativePlatform()`,所以共用 Native 初始化代码 | 存在 Xcode/Capacitor 工程;但多个 UI/runtime 分支只检查 Android,不能据此宣称 iOS 等价;Windows 无 iOS 构建/设备证据 | +| HarmonyOS | ArkTS/ArkWeb + `nowen-harmony/entry/src/main/ets/bridge/JsBridge.ets` | 没有检索到共用 Native SQLite Sync V2 的桥接实现;不能宣称离线强杀等价于 Android | + +既有并行链路:Sync V2、Web IndexedDB/REST 离线队列、旧 `/offline-sync` feed、Yjs 文档协作、Electron 文件夹同步。Android 的 `OfflineSyncRuntime` 已明确不挂载旧 Web 队列 runtime,但知识树/正文仍跨两种权威数据源。Yjs 是打开文档的协作能力,不是完整设备持久化协议;本次不引入第三套同步协议。 + +## 功能与实体覆盖矩阵(源码证据) + +“实现”不代表跨平台生产验收已完成。 + +| 能力/实体 | Desktop V2 | Android V2 | Web | 证据/缺口 | +| --- | --- | --- | --- | --- | +| Markdown/富文本笔记、标题、常用元数据 | 实现 | 实现 | REST + cache/queue | `apply.ts`/`applyLocal.ts`/`mobileSyncEngine.ts`;正文按完整文档版本处理,不拼接富文本 JSON | +| 笔记本、标签、note_tag、收藏 | 实现 | 实现 | REST | 部分实体没有正文级版本冲突协议,不能承诺任意并发字段修改都可恢复 | +| knowledge_tree_node、结构和排序 | 协商能力,已有 baseline/readiness/冲突路径 | **未接入**移动 V2 实体注册表;登录树走服务器 | `types.ts` 额外协商树实体;`isCoreEntityType()` 只有十类;Native 纯设备树为投影 | +| 任务、task_reminder、日记、思维导图 | 实现 | 实现 | REST 与既有模块离线适配 | 高级任务项目/模板/依赖和附件文件夹不在十类实体集合 | +| 附件元数据和二进制 | 分开追踪 | 分开追踪 | 缓存/恢复路径 | `blob.ts`/`sync-v2-blob.ts`/`transferAttachments()`;笔记正文 ACK 不等于附件成功 | +| 回收站/恢复/永久删除 | 已有路径 | 已有路径,但删除保护存在 P0 缺口 | REST + cache | 历史清理后的 tombstone/reset 和父实体级联仍需加强回归 | +| 加密笔记 | 密文 guard | 密文 guard | 密文 REST + client decrypt | `encrypted-notes-sync.test.ts`、`noteDocument.ts`;不得把正文/密钥写日志 | +| workspace、权限撤销 | scope/fingerprint/冻结 | scope/fingerprint/冻结 | REST ACL | 必须保留跨用户/撤权负例;离线密码目录另有访问 guard | +| 本地保存与 Outbox 原子性 | SQLite 触发器事务 | SQLite transaction | **未完成统一原子链路** | 浏览器不能被描述为完整 Local-first | +| 持久化 mutationId、重复 Push | 实现 | 实现 | REST 回放并非所有入口都有稳定 mutationId | `sync_applied_mutations`;重试不得换 mutation 身份 | +| 逐笔记 ACK 收据 | 本机 REST 不当云端确认,查询不完整 | PR #821 mutation ACK | 本次修订 + 有效正文版本 ACK | 没有可靠另一设备收据,不能称所有设备已同步 | +| 通知/重连/周期 Pull | backend realtime + engine timer | 缺少 runtime `sync.changed` → Pull 接线与稳定周期补偿 | 既有 realtime + REST refresh | 前台通知性能尚无 P50/P95/P99 | +| 冲突查看/保留两端/另存为 | 已有设置 UI/冲突记录 | 已有 admin adapter/冲突记录 | queue/draft/conflict store | 删除冲突和解决后的版本基线需专项回归 | + +## P0/P1/P2 问题清单 + +| 级别 | 代码证据与影响 | 当前复现状态/后续验证 | +| --- | --- | --- | +| P0 | Native Push ACK 仅保存 receipt/删除队列,未校准 `notes.version` 和后续队列 `baseVersion`;桌面 `acknowledgeNotePush()` 已做校准 | 需独立 SQLite/真实 API 回归确认连续编辑及 ACK 延迟,避免自身假冲突 | +| P0 | Native `applyEntries()` 的 pending/conflict 判断排除了 `__delete`;删除可直接覆盖本地尚未上传正文 | 需失败用例证明本地正文/Outbox/双方冲突在远端删除后仍可恢复 | +| P0 | Native `request()` 在响应 JSON 异步解析后没有再次校验停止/登录身份;旧会话响应可能继续 Apply | 需延迟 JSON、切账号/stop 回归,旧 DB 不允许接受迟到数据 | +| P0 | Native ACK 收据允许 mutationId 存在但 version 非法仍显示 confirmed | 需非法/缺失版本用例;无证据应为 unverified | +| P1 | 登录态树创建与正文数据库不同,旧版可见树节点但读本地不存在 | #829 原文无日志;当前基准已有授权远端 fallback 和回归,不能宣称复现或真机已解决;fallback 仍非离线本地可读 | +| P1 | Native runtime 没有订阅 `sync.changed`、没有普遍周期 Pull;收到事件与 UI 刷新也需以本地 Apply 为界 | 需接线测试、通知丢失/重连及实际延迟测量 | +| P1 | pull 的同实体多次 upsert/delete 未按最终 feed 操作折叠;snapshot 可能缺失已删除实体而一直报错 | 需 create→delete 同页、delete→recreate 同页回归;snapshot 与 feed 竞态需明确协议边界 | +| P1 | Snapshot pruning 使用父实体删除/附件物理删除,保护仅检查该实体的 Outbox | 需父子依赖、冲突、附件队列保留检查,不能通过删库修复 | +| P1 | 冲突解决 Native keep-local 仍使用本地 payload 版本,服务端版本更高时后续编辑可能重新冲突 | 需双方保留/版本收敛/删除恢复测试 | +| P2 | 文档 `sync-realtime.md` 主要描述旧 Web cache/full-snapshot 路径 | 需更新平台边界与当前 V2 行为 | + +Android“Web 修改后闪退”尚未取得本次复现/新崩溃日志,不能用 TypeScript 异常处理替代原生诊断。既有 Kotlin Filesystem 包装修复和 Native list 预览限长在基准中;本次必须检查 logcat 和隔离包运行证据。 + +## 测试缺口与安全策略 + +已有大量单元/路由测试,但 `sync-v2-engine.test.ts` 使用可编程 FakeRemote;不能替代三个独立客户端数据库与一个真实服务端数据库。新增测试应使用真实路由、独立文件和持久化重开,覆盖双向/三端、并发、ACK 后崩溃、Apply 失败、删除、非法 ACK、旧会话回调和权限拒绝。故障日志只保留错误码/耗时/计数,不输出正文、Token 或密钥。 + +禁止清库/清 Outbox 解决普通故障。先在测试库证明原子回滚和重启恢复,再修改业务代码;不需要数据迁移的修复保持 schema 不变。涉及设备验收使用独立包与测试账号,不触碰现有用户笔记。Docker、NAS、macOS/Linux 包、iOS 和 HarmonyOS 没有真实执行证据时均列为未验收。 + +实施顺序与成功标准: + +1. Phase 0 审计并保留证据 → 本文可追溯到实际调用链、Issue #829 和 PR #821。 +2. Phase 1 先写失败测试再修 P0 → 本地正文/队列不丢、ACK 仅确认发送的版本、旧会话不能 Apply、双方冲突可恢复。 +3. Phase 2 接通知、周期补偿与 UI → 通知只唤醒 Pull,Apply 后刷新,无回声循环,记录实测延迟。 +4. Phase 3 回归与交付 → 独立提交、目标构建/lint/设备证据、PR 指向 `release/v1.5.2`,剩余项明示;不发布正式版本。 diff --git a/docs/sync-v2-reliability-validation.md b/docs/sync-v2-reliability-validation.md new file mode 100644 index 000000000..cc9021d2f --- /dev/null +++ b/docs/sync-v2-reliability-validation.md @@ -0,0 +1,121 @@ +# Sync V2 reliability validation / 同步可靠性验证 + +日期:2026-10-11。基准 `d95d3787`;开发分支 `codex/sync-v2-reliability`。所有笔记正文均为临时测试数据。 + +## Phase 0 + +审计交付:`docs/sync-v2-reliability-audit.md`,commit `bd26fabe`。已核查 Issue #829 和合并后的 PR #821。架构图、平台读写路径、实体矩阵、问题清单和缺口见审计文件。 + +## Phase 1:数据一致性修复 + +| 已复现问题 | 修复与执行链路 | +| --- | --- | +| 远端永久删除丢掉本机离线正文 | `mobileSyncEngine.ts`:删除也经过 pending/冲突保护;正文、mutation 与远端删除标记同时保留 | +| 旧离线更新复活已永久删除的服务器笔记 | `apply.ts`/`sync-v2.ts`:带已有基线的更新不再当创建;返回可恢复删除冲突;恢复须另存新身份 | +| 本机同时保存标题/正文丢掉先前修改 | `nativeLocalRepository.ts`:读取当前笔记/权限/版本、更新业务和插入 Outbox 同个串行事务 | +| 无效 ACK 显示云端已确认 | `mobileNoteSyncReceipt.ts`:必须有有效正整数版本;`mobileSyncEngine.ts` 在出队前验证当前批次和版本,非法 ACK 保留队列 | +| 在途切账号后 JSON 解码继续接受旧 ACK | `mobileSyncEngine.ts`:传输取消、重试发送前校验 AbortSignal,JSON 解码后再检查会话;旧会话不宣告状态变化 | +| ACK 等待期间继续编辑导致自身冲突 | ACK 事务校准后续 baseVersion;批次内同笔记按服务端应用顺序接续版本;自身回声只认精确 mutation ACK 且读取最新 pending payload | +| 未记录本地 ACK 时崩溃,重启 Snapshot 误判已创建实体为冲突 | 持久化 mutationId 先重放 Push,再做首次/恢复 Snapshot;duplicate 结果与原队列出队同事务落盘 | +| 解决冲突后版本基线分叉 | keep-local 以远端版本生成新修订;冲突查询限定 Profile;未解决冲突阻止同实体的新队列静默上传 | +| Desktop 合并离线创建和编辑后携带错误已有版本 | `push.ts`:保持最早的缺失 baseVersion,创建不会因后续编辑变成更新 | + +未使用清库/删除全部 Outbox/覆盖真实用户笔记;没有 schema 迁移或新增依赖。冲突台账保留双方。拒绝旧身份更新会让过去悄悄“重新创建”的客户端产生明确冲突;这是预期保护。永久删除后的保留本地须另存为新笔记,不能复用旧 ID。 + +自动化结果: + +- `sync-v2-mobile-multidevice.test.ts`:12/12。独立磁盘上的三客户端 SQLite + 独立服务器 SQLite,真实 HTTP/Hono Sync V2 路由;覆盖 Markdown/富文本、双向和三端收敛、并发、close/reopen、丢失 ACK、重复 mutationId、Apply/Outbox 失败、永久删除、自身回声、非法收据。 +- 后端组合 `sync-v2-protocol.test.ts`、`sync-v2-engine.test.ts`、`sync-v2-mobile-multidevice.test.ts`、`encrypted-notes-sync.test.ts`:84/84。 +- 前端目标组合 `mobileSyncEndpoint`、`mobileNativeReceiptIntegration`、`mobileSyncUnknownEntity`、`mobileSyncUnexpectedFailure`、`mobileNoteSyncReceipt`、`nativeLocalNoteList`:35/35;另一次含目录密码与 #829 既有 fallback 的组合:70/70。 +- `node scripts/data-consistency-contract.mjs --verify`:12 条契约、56 个用例引用;新增 DCC-012 已进入既有门禁。 +- 前端 `npm run build`:通过。此后源码变更仍须最终再构建;不能用该中间构建证明最终 APK。 + +测试执行方式:Windows 的 `better-sqlite3` 是 Electron ABI,必须使用 `ELECTRON_RUN_AS_NODE=1` 的项目 Electron,参数为 `--import ./backend/node_modules/tsx/dist/loader.mjs --import ./backend/tests/setup-db-isolation.ts --test <目标文件>`。PowerShell 对 GUI exe 应使用 `Start-Process -WindowStyle Hidden -Wait -PassThru` 并检查 `ExitCode`。直接 Node 因 ABI 不符没有运行成功;未加载隔离 setup 的加密套件最初失败,补正确 setup 后已通过。新增测试最初确实出现删除丢正文、错误 confirmed、复活旧笔记、并行回写丢标题和延迟 ACK 假冲突的失败。 + +## Phase 2:通知、界面与恢复 + +| 根因 | 改动与新链路 | +| --- | --- | +| Native 无 `sync.changed` → Pull 接线 | `mobileSyncRealtime.ts`/runtime:复用既有 realtime;通知/重连/前台 → 150ms 合并唤醒 → 同一个引擎 Pull;30s 前台补偿 | +| 原生 Apply 后列表/打开正文没有明确提交事件 | `mobileSyncEngine` 提交后事件 → `OfflineSyncRuntime` 的 Native 刷新分支与 `EditorPane` 既有版本/未保存内容保护;旧引擎回调被当前 runtime 身份挡住 | +| 创建与永久删除在同 Feed 页使 Snapshot 缺项 | Native/Desktop 均按身份折叠最终操作;不跳过真正缺失的 upsert,不推进失败游标 | +| Native 全局事务标记误伤独立 UI 查询 | `nativeDatabase.ts`:外层请求进入既有串行队列等待提交/回滚;事务回调仍只能用 tx;保留失效 tx、嵌套事务、关闭连接与加密 DDL 检查 | +| Workspace Snapshot 清理级联删待同步子项/冲突/附件 | Native/Desktop 排除 unresolved 冲突并保护剩余父子关系;Native 保留未上传本地附件,只有元数据确实移除才删除文件 | +| 冲突后继续编辑,keep-local/fork 使用过期快照 | Native 选择时在事务内读最新本地实体,随后生成远端基线的新修订;原冲突两方台账不删除 | + +对应失败测试先失败再修复;多库套件增至 14 项。前端包括 Native UI 刷新、通知调度、冲突“查看”/知识树安全选择、目录密码保护与 401 endpoint 恢复。Node 24 的真实 SQLite 本地账号导入/同步开关组合 16/16,覆盖附件复制失败可重试与模式切换不丢队列。旧同步开关 fixture 缺少成功 ACK `version`,已按真实 API 契约修正;没有放宽生产校验。 + +没有引入第三套同步协议、schema 迁移或新增依赖。父实体因受保护子项被保留是保守的数据保护,之后仍需用户处理冲突。Native 外层数据库队列不能在其自身事务回调内被 await,须使用 tx;本次修改及迁移调用链遵守该约束。 + +## 未验收与下一步 + +- #829 基准已有“本地缺失 → 服务端授权 → 远端完整读取”fallback;本次没有新复现日志。该 fallback 仍不等于刚创建的服务端树笔记已经持久化到 Native,不能宣称离线完整解决。 +- Android 原生闪退:已读取连接设备 PJZ110 crash buffer,没有获得本次 Sync V2 的明确崩溃堆栈;原生列表预览限长与 Kotlin Filesystem 修复在基准里。新增源码验证不能替代新包设备验收。 +- 知识树 Native V2 实体未实现,登录态仍服务端树;附件 bytes、完整 EditorPane 与树创建仍需新包专项验收。父实体增量删除与所有受保护后代的组合须继续扩展故障覆盖,不能用本次 Workspace Snapshot 回归替代。 +- Phase 3 的最终构建、设备实测与完成度矩阵见后续记录;Docker/NAS、iOS/HarmonyOS/macOS/Linux 安装包仍需独立平台证据。 + +## Phase 3:本轮回归与设备证据(整体生产验收未完成) + +生产源码提交:Phase 0 `bd26fabe`;Phase 1 `421d37eb`;Phase 2 `4588de41`。均在 `codex/sync-v2-reliability`,目标 `release/v1.5.2`。Phase 3 单独提交验收工具和记录。保留 Draft PR,不合并、不发布。 + +| 检查 | 最终结果与边界 | +| --- | --- | +| backend 全部 `sync-v2-*.test.ts` + encrypted sync + backup integrity + existing-tree migration + native encrypted storage | **272/272**;Electron RunAsNode、隔离 SQLite,串行 test runner。新增多库测试 **14/14** | +| frontend 15 个目标文件 | **121/121**;Node 24 + Vitest;真实 SQLite migration/sync-switch 套件没有 skip;含音频 ACK 兼容 | +| consistency gate | **12 invariants / 59 references**;新增 Native 事务并发、删除折叠、Snapshot 保留均在门禁 | +| frontend `npm run build`、backend `npm run build` | 通过;已有大 chunk 警告仍在 | +| frontend `npm run lint:baseline` | **1447 historical errors / 0 new errors**;不声称全仓无 lint 历史错误 | +| Android `assembleDebug` | 通过;JDK 21、API 35 emulator,独立包名 `com.nowen.note.syncacceptance`;不覆盖正式包 | +| Android SQLite / 强杀 / 同步 | 真正 Native 插件;移除 adb reverse 验证服务器不可达 → MD/富文本创建立即读取 → `am force-stop` → 离线重开同 ID/正文/Outbox → 恢复网络自动上传/ACK;通过 | +| Android 通知/恢复 | HTTP 客户端 Push → 实际 WS → Native Pull;双向、主动断开通知后的 30s 补偿、重连;通过。测试期间该包 PID crash buffer 0 行 | +| Docker/NAS、其他系统安装包 | 本地 **未运行**;本机无 Docker CLI/runtime;远端生产 image/plugin smoke 见下段;没有 NAS/其他平台设备证据 | + +PR #830 的 `6b4c5949` 远端 Data Consistency Contract CI(三个 job)、Android Local-first、Knowledge Tree、Data Protection CI 已通过;Docker Plugin Artifacts CI 已真正构建生产镜像并在镜像内执行 plugins smoke,通过,但不等于 NAS/完整同步镜像验收。Voice Memo 最初有 3 条失败:旧 fixture 成功 ACK 缺少 `version`、未模拟后续 Outbox baseVersion 更新,按真实协议补齐后本地 3/3。I18n Release CI 的 16 个 `encryptedConversion.*` 缺失键在原始 `d95d3787` 再跑也失败(23 pass / 1 fail),属于基准问题,未混入同步修复。后续提交的 CI 须以对应 SHA 查询,不能把旧 SHA 的通过状态套用到新提交。 + +Android 源码来自 `4588de41`,测试工具另在本阶段。隔离 APK SHA256:`4399aefd46c818de3eba68f8c19e7d2e77887de51a14cb572b1c79535b7f9622`。模拟器 `Pixel_8_API_35`,API 35,WebView 124.0.6367.219;只读启动现有 AVD,不改正式 App 数据。真实连接的 PJZ110 上安装未返回,本次没有真机验收通过证据。 + +最终机器记录:[Android report](test-results/sync-v2-android-2026-10-11.json)。各方向 30 次、同一 host 的 adb reverse 回环网络,计时包括轮询观察开销: + +| 链路 | P50 | P95 | P99(30 样本取最大值) | +| --- | ---: | ---: | ---: | +| HTTP Push 发出 → Android Native Pull 提交后读到正文 | 337ms | 348ms | 768ms | +| Android edit 发出 → 服务端正文 + Native ACK 落盘 | 375ms | 394ms | 429ms | + +这不是公网 P95/P99,也不是 PC GUI 或完整 EditorPane 的端到端指标;样本数不足以估计生产尾延迟。目标 P95≤3s/P99≤10s 仍须实际部署网络验收。WebSocket 只传唤醒通知,测量没有主动调用每次 sync 来替代自动链路。 + +失败记录:最初设备脚本通过 Chromium CDP 连接不兼容 WebView context 管理,改用 Playwright Android API;localStorage 测试 ID 标记强杀未及时落盘,改为 SQLite 测试元数据;真实 Native 并发查询错误已补失败测试并修复;超过 15 分钟的旧 fixture JWT 导致 ACK 等待失败,队列保留,修正测试工具每轮签发后通过。没有对生产 Token/ACK 校验放宽。日志保留在工作树忽略的 `sync-*.log`,只含临时测试数据;提交的 report 没有 Token、正文或密钥。 + +### 完成度矩阵(代码支持不等于平台验收) + +| 能力/实体 | Desktop | Android Native | Web/移动浏览器 | 本轮证据/剩余工作 | +| --- | --- | --- | --- | --- | +| 正文、标题、格式、元数据;MD/富文本 | V2 + SQLite | V2 + SQLite | REST + cache/queue | 独立多库双向/三端;Native 设备两格式;完整 GUI 跨端未验收 | +| 本地保存 + Outbox 原子性 | 事务触发器 | 同事务读/写/Outbox | **统一原子 Local-first 未实现** | 失败入队回滚、并发保存;浏览器仍是部分离线链路 | +| mutationId 幂等 / ACK / 自身回声 | 实现 | 修复并回归 | 分入口实现 | 丢 ACK 重放、连续编辑、非法 ACK;另一设备收据未实现 | +| 自动实时 / 重连 / 前台补偿 | 既有 backend engine | 新接线,设备通过 | 既有 realtime/queue | Android 30s 漏通知补偿;系统强杀后后台实时不支持 | +| 离线重启恢复 | 磁盘 SQLite | 实际强杀通过 | 缓存/草稿,部分 | Native 新建/编辑持久化;PC 强杀与浏览器事务故障专项未运行 | +| 并发冲突、查看、选择、另存为 | 已有台账/UI | 最新编辑/版本基线修复 | 已有 draft/conflict | 双方保留/keep-local/fork;树结构选择现有测试通过;完整手机 UI 未验收 | +| 笔记本、文件夹、排序 | legacy V2 | legacy V2 | REST | 代码/协议回归;不是完整知识树离线能力 | +| 知识树节点、移动、结构冲突 | 协商的 V2 实体 | **Native V2 节点未实现**,登录态 server-first | REST | #829 已有授权 fallback;树创建后离线立即打开仍缺口 | +| 标签、收藏、关联 | V2 | V2 | REST/queue | 后端实体回归;Native 设备关联专项未运行 | +| 任务、状态、提醒 | legacy V2 | legacy V2 | REST | 代码/协议回归;Native 设备专项未运行 | +| 日记、思维导图 | legacy V2 | legacy V2 | REST | 同上;新模型/复合资源不冒充全覆盖 | +| 附件元数据 | V2 | V2 | REST/cache | 当前 metadata 与 bytes 分开;受保护元数据/本地文件清理测试通过 | +| 附件 bytes/hash/失败恢复 | 独立 transfer | 独立 transfer | offline attachment jobs | 路由/迁移已有回归;跨端新包 bytes 设备专项未运行,不宣称整篇附件完成 | +| 回收站、恢复、永久删除 | 已有 + 版本保护 | note fields/delete + 本地保护 | REST | 永久删除/旧更新不复活通过;回收站往返设备专项未运行 | +| Tombstone / Snapshot 历史清理 | reset/Snapshot | reset/Snapshot,保留冲突/子项 | 旧 snapshot/cache | 本次同页折叠与 pruning 回归;无基线旧 ID 的长期 tombstone 策略仍需评估 | +| 账号 / Profile / 服务器隔离 | Profile/scope | 旧会话取消/响应检查 | scoped cache | 延迟切账号负例;完整双服务器 GUI 切换未验收;不清未同步数据 | +| Workspace / ACL / 撤权 | scope/fingerprint | plan/freeze/guard | server ACL | 后端权限负例;设备完整撤权/密码目录验收仍需补齐 | +| V2 加密正文安全 | 密文验证/guard | guard + 密文同步 | 客户端加密 | encrypted sync/native storage 回归通过;加密+附件设备专项未运行 | +| iOS / HarmonyOS | — | iOS 部分共享;Harmony 缺 Native V2 bridge | — | **均未等价验收**;不能套用 Android 结论 | + +后续阻断生产验收的工作:Web 统一 IndexedDB 事务/Local-first;Native 知识树离线创建及 #829 原路径;父实体增量删除与所有受保护后代组合;无基线旧 ID tombstone;完整移动 EditorPane/冲突/附件/权限设备场景;公网长时稳定性与真实 P99;Docker/NAS/macOS/Linux/iOS/HarmonyOS 验收。不能将本轮修复或 CI 当作这些项已经完成。 + +### 重跑隔离 Android 验收 + +入口只由 `vite.android-sync.config.ts` 打进测试包,正常 `npm run build` 不包含测试 HTML/工具接口。后端 helper 只监听 127.0.0.1,不部署到生产。 + +1. 根目录:配置临时 `SYNC_ACCEPTANCE_DIRECTORY`,用项目 Electron RunAsNode + tsx 启动 `backend/tests/sync-v2-android-server.ts`;端口 47831。不要使用生产 DB。 +2. frontend:`npx vite build --config vite.android-sync.config.ts`;Node≥22 运行 Capacitor `sync android`;JDK21 执行 `gradlew -I :app:assembleDebug`。 +3. 只向隔离设备安装 APK,`adb -s reverse tcp:47831 tcp:47831`,启动 `com.nowen.note.syncacceptance/com.nowen.note.MainActivity`。 +4. 根目录:`node scripts/android-sync-acceptance.mjs `。只操作这个测试包;关闭其网络转发来模拟离线,恢复后测量。report 写到 `docs/test-results/`。结束后关闭 helper、移除转发并关闭测试模拟器。 diff --git a/docs/sync-v2-reliability.en.md b/docs/sync-v2-reliability.en.md new file mode 100644 index 000000000..a72e3dd02 --- /dev/null +++ b/docs/sync-v2-reliability.en.md @@ -0,0 +1,30 @@ +# Sync V2 reliability: implementation and acceptance + +This branch improves the existing Sync V2 protocol; it does not replace it with another synchronization system. It is a **Draft candidate**, not completed production acceptance. Base: `release/v1.5.2` at `d95d3787`; branch: `codex/sync-v2-reliability`. + +The [audit](sync-v2-reliability-audit.md) traces each platform's actual storage and synchronization path. The [validation record](sync-v2-reliability-validation.md) contains the complete entity/platform matrix, failure evidence, commands and limitations. + +## Implemented fixes + +- Native note writes read the latest version, save content and enqueue mutations in one transaction. Independent UI reads wait for Apply to finish instead of failing on a global transaction flag. +- Successful ACKs must identify a sent mutation and contain a valid server revision. Pending descendants rebase on ACK; old account responses are rejected after asynchronous JSON decoding. Restart retries the durable mutation before reconciling Snapshot. +- Remote deletion preserves unsent local content as a conflict. An update with an existing base cannot silently recreate a permanently deleted server note. Restoring that content uses a new identity. +- A native conflict choice/fork reads the latest durable local edit. Workspace Snapshot pruning protects unresolved conflicts, pending descendants and local attachment transfers; file removal follows actual metadata removal. +- Native runtime reuses `/ws` notices/reconnect to wake Pull. Foreground recovery and a 30-second foreground poll compensate for missing notices. UI refresh follows committed SQLite Apply. WebSocket carries notices, never authoritative note bodies. +- Desktop and Native fold multiple operations per identity in a Change Feed page before requesting current Snapshot payloads. + +## Evidence + +Final targeted runs: backend **272/272**, frontend **121/121**, consistency gate **12 invariants / 59 references**. Frontend/backend builds pass; lint baseline reports **0 new errors** with 1447 historical errors. No schema migration or dependency was added. + +Remote checks at `6b4c5949`: Data Consistency, Android Local-first, Knowledge Tree and Data Protection CI pass. Docker CI builds the actual production image and runs the plugin smoke inside it successfully. Voice memo fixtures initially omitted successful ACK revisions; updating them to the real protocol yields 3/3 locally. The 16 missing encrypted-conversion translation keys also fail on base `d95d3787`, so that I18n failure is recorded as baseline debt. New commits require their own CI results. + +The isolated Android API 35 emulator package used real Native SQLite. With the test HTTP endpoint unreachable, Markdown/rich-text creation was immediately readable. Force-stop/relaunch preserved identity, content and pending mutations; reconnection uploaded them and recorded ACKs. Actual WS/Pull, missed-notice polling and reconnect passed. The test package's crash buffer was empty during this run. + +In an adb loopback environment, 30 samples per direction measured P50/P95/P99 of **337/348/768ms** for HTTP Push to Native Pull, and **375/394/429ms** for Native editing to server content plus persisted ACK. These are narrow integration measurements, not production WAN or full editor GUI latency. See the [machine report](test-results/sync-v2-android-2026-10-11.json). + +## Remaining acceptance blockers + +Web/mobile browsers remain REST plus IndexedDB cache/offline queue, without unified transactional Local-first storage. Signed-in Android knowledge tree operations remain server-first; Native knowledge-tree V2 is missing. Issue #829's existing authorized remote fallback is not proof that tree-created notes are durably readable offline. Full EditorPane, tree creation, conflict UX, attachment bytes and revocation need device acceptance. Parent deletion with protected descendants, no-base stale IDs/tombstones, long-running WAN latency and crash reproduction need further tests. Docker full-sync/NAS and macOS/Linux/iOS/HarmonyOS package acceptance have not run. + +Commits: audit `bd26fabe`, durability fixes `421d37eb`, native realtime/recovery `4588de41`. No main modification, force-push, merge or release was performed. The [Android harness](../scripts/android-sync-acceptance.mjs), test-only Vite configuration and loopback backend helper are excluded from normal production builds. diff --git a/docs/test-results/sync-v2-android-2026-10-11.json b/docs/test-results/sync-v2-android-2026-10-11.json new file mode 100644 index 000000000..596272213 --- /dev/null +++ b/docs/test-results/sync-v2-android-2026-10-11.json @@ -0,0 +1,28 @@ +{ + "device": "emulator-5580", + "package": "com.nowen.note.syncacceptance", + "offlineImmediateRead": true, + "forceStopRecovery": true, + "offlineNetworkIsolation": "adb reverse removed; HTTP server unreachable from Android", + "markdownAndRichText": true, + "missingNoticePoll": true, + "reconnect": true, + "httpPushToNativePull": { + "samples": 30, + "p50Ms": 337, + "p95Ms": 348, + "p99Ms": 768 + }, + "nativeCommitToServerAndAck": { + "samples": 30, + "p50Ms": 375, + "p95Ms": 394, + "p99Ms": 429 + }, + "errors": [ + { + "lastError": null + } + ], + "boundary": "Loopback adb reverse; real Android native SQLite and MobileSyncEngine. Does not test full EditorPane, knowledge tree creation, attachment bytes, or production JWT middleware." +} \ No newline at end of file diff --git a/frontend/src/components/EditorPane.tsx b/frontend/src/components/EditorPane.tsx index d46767eea..2811bb075 100644 --- a/frontend/src/components/EditorPane.tsx +++ b/frontend/src/components/EditorPane.tsx @@ -25,6 +25,7 @@ import { parseMermaidMindmap, normalizeMindMapData } from "@/lib/mindmapTransfor import { cn } from "@/lib/utils"; import SyncStatusBadge from "@/components/SyncStatusBadge"; import { isAndroidNativeRuntime } from "@/lib/mobileLocalMode"; +import { MOBILE_SYNC_APPLIED_EVENT, type MobileSyncAppliedDetail } from "@/lib/mobileSyncStatus"; import { applyEditorUpdateToNote, PREPARE_EDITOR_SPLIT_CLOSE_EVENT, @@ -1281,6 +1282,20 @@ function OrdinaryEditorPane({ } } + useEffect(() => { + const applied = (event: Event) => { + const cur = activeNoteRef.current; + const detail = (event as CustomEvent).detail; + if (!cur || !detail) return; + if (detail.deletedNoteIds.includes(cur.id)) setRemoteDelete({ trashed: false }); + else if (detail.noteIds.includes(cur.id)) void checkActiveNoteRemoteVersion("native-sync"); + }; + window.addEventListener(MOBILE_SYNC_APPLIED_EVENT, applied); + return () => window.removeEventListener(MOBILE_SYNC_APPLIED_EVENT, applied); + // The handler reads current note/editor state through existing refs. + // eslint-disable-next-line react-hooks/exhaustive-deps + }, []); + const { presenceUsers, isConnected, setEditing: rtSetEditing } = useRealtimeNote({ noteId: activeNote?.id ?? null, // ��ʽ���� selfUserId��EditorPane ������ selfUser��localStorage ���� + /api/me���� diff --git a/frontend/src/components/OfflineSyncRuntime.tsx b/frontend/src/components/OfflineSyncRuntime.tsx index f2bd2c246..ac1369ee6 100644 --- a/frontend/src/components/OfflineSyncRuntime.tsx +++ b/frontend/src/components/OfflineSyncRuntime.tsx @@ -5,6 +5,7 @@ import { installOfflineAttachmentRecoveryCapture } from "@/lib/offlineAttachment import { useNetworkStatus } from "@/hooks/useNetworkStatus"; import { useAppActions } from "@/store/AppContext"; import { isAndroidNativeRuntime } from "@/lib/mobileLocalMode"; +import { MOBILE_SYNC_APPLIED_EVENT } from "@/lib/mobileSyncStatus"; /** * Web / Electron 使用的旧离线队列 Runtime。 @@ -36,6 +37,22 @@ function ServerOfflineSyncRuntime() { } export default function OfflineSyncRuntime() { - if (isAndroidNativeRuntime()) return null; + if (isAndroidNativeRuntime()) return ; return ; } + +function NativeSyncRefresh() { + const actions = useAppActions(); + useEffect(() => { + const refresh = () => { + actions.refreshNotes(); actions.refreshNotebooks(); + window.dispatchEvent(new CustomEvent("nowen:knowledge-tree-changed", { detail: { reason: "native-sync" } })); + void api.getTags().then(actions.setTags).catch((error) => { + console.warn("[OfflineSyncRuntime] refresh tags after native sync failed:", error); + }); + }; + window.addEventListener(MOBILE_SYNC_APPLIED_EVENT, refresh); + return () => window.removeEventListener(MOBILE_SYNC_APPLIED_EVENT, refresh); + }, [actions]); + return null; +} diff --git a/frontend/src/components/__tests__/NativeSyncRefresh.test.tsx b/frontend/src/components/__tests__/NativeSyncRefresh.test.tsx new file mode 100644 index 000000000..15c51a335 --- /dev/null +++ b/frontend/src/components/__tests__/NativeSyncRefresh.test.tsx @@ -0,0 +1,36 @@ +import { act } from "react"; +import { createRoot } from "react-dom/client"; +import { expect, it, vi } from "vitest"; +import OfflineSyncRuntime from "../OfflineSyncRuntime"; +import { MOBILE_SYNC_APPLIED_EVENT } from "@/lib/mobileSyncStatus"; + +const actions = vi.hoisted(() => ({ refreshNotes: vi.fn(), refreshNotebooks: vi.fn(), setTags: vi.fn() })); +const legacy = vi.hoisted(() => ({ useNetworkStatus: vi.fn(), syncAttachments: vi.fn() })); +vi.mock("@/store/AppContext", () => ({ useAppActions: () => actions })); +vi.mock("@/lib/api", () => ({ api: { getTags: async () => [] } })); +vi.mock("@/lib/mobileLocalMode", () => ({ isAndroidNativeRuntime: () => true })); +vi.mock("@/hooks/useNetworkStatus", () => ({ useNetworkStatus: legacy.useNetworkStatus })); +vi.mock("@/lib/offlineAttachmentRecovery", () => ({ installOfflineAttachmentRecoveryCapture: legacy.syncAttachments })); +vi.mock("@/lib/syncEngine", () => ({ SYNC_SNAPSHOT_APPLIED_EVENT: "legacy-snapshot" })); + +it("refreshes native lists and tree after Apply without starting a second sync engine", async () => { + Object.assign(globalThis, { IS_REACT_ACT_ENVIRONMENT: true }); + const host = document.createElement("div"); + const root = createRoot(host); + const treeChanged = vi.fn(); + window.addEventListener("nowen:knowledge-tree-changed", treeChanged); + try { + await act(async () => root.render()); + expect(legacy.useNetworkStatus).not.toHaveBeenCalled(); + expect(legacy.syncAttachments).not.toHaveBeenCalled(); + expect(actions.refreshNotes).not.toHaveBeenCalled(); + await act(async () => { window.dispatchEvent(new CustomEvent(MOBILE_SYNC_APPLIED_EVENT)); }); + expect(actions.refreshNotes).toHaveBeenCalledOnce(); + expect(actions.refreshNotebooks).toHaveBeenCalledOnce(); + expect(actions.setTags).toHaveBeenCalledWith([]); + expect(treeChanged).toHaveBeenCalledOnce(); + await act(async () => root.unmount()); + window.dispatchEvent(new CustomEvent(MOBILE_SYNC_APPLIED_EVENT)); + expect(actions.refreshNotes).toHaveBeenCalledOnce(); + } finally { window.removeEventListener("nowen:knowledge-tree-changed", treeChanged); } +}); diff --git a/frontend/src/lib/__tests__/mobileNativeReceiptIntegration.test.ts b/frontend/src/lib/__tests__/mobileNativeReceiptIntegration.test.ts index e1a3f7f4e..b9ef9d660 100644 --- a/frontend/src/lib/__tests__/mobileNativeReceiptIntegration.test.ts +++ b/frontend/src/lib/__tests__/mobileNativeReceiptIntegration.test.ts @@ -1,5 +1,5 @@ // @vitest-environment jsdom -import { afterEach, describe, expect, it, vi } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { MobileSyncEngine, isAckedOwnNoteEcho } from "../mobileSyncEngine"; import { nativeReceiptKey } from "../mobileNoteSyncReceipt"; import type { NativeDatabase } from "../nativeDatabase"; @@ -12,6 +12,7 @@ const local = { id: "note-a", scopeKey: "personal", version: 4, title: "RACE NEW const ownAck = JSON.stringify({ mutationId: "mutation-old", serverVersion: 3 }); const pending = { baseVersion: 3, payload: JSON.stringify(local) }; +beforeEach(() => { localStorage.clear(); }); afterEach(() => { vi.unstubAllGlobals(); }); describe("PR #821 Android receipt ACK and in-flight pull regressions", () => { diff --git a/frontend/src/lib/__tests__/mobileSyncEndpoint.test.ts b/frontend/src/lib/__tests__/mobileSyncEndpoint.test.ts index 70b93b2b1..13546f508 100644 --- a/frontend/src/lib/__tests__/mobileSyncEndpoint.test.ts +++ b/frontend/src/lib/__tests__/mobileSyncEndpoint.test.ts @@ -66,6 +66,38 @@ describe("mobile sync transport rerouting", () => { expect(db.transaction).not.toHaveBeenCalled(); }); + it("rejects JSON that completes after its sync session was stopped", async () => { + vi.useFakeTimers(); + const { engine, db } = createEngine(); + let finish!: (value: unknown) => void; + const body = new Promise((resolve) => { finish = resolve; }); + const json = vi.fn(() => body); + vi.stubGlobal("fetch", vi.fn(async () => ({ ok: true, json }))); + engine.start(); + const pending = engine.syncOnce(); + while (!json.mock.calls.length) await Promise.resolve(); + engine.stop(); + finish({ items: [ { scopeKey: "personal", workspaceId: null, canWrite: true, accessFingerprint: "old" } ] }); + await pending; + expect(db.transaction).not.toHaveBeenCalled(); + expect(db.run).not.toHaveBeenCalled(); + }); + + it("rejects a late push ACK when the account changes while decoding JSON", async () => { + localStorage.setItem("nowen-token", "same-token"); + const { internal } = createEngine(); + let finish!: (value: unknown) => void; + const body = new Promise((resolve) => { finish = resolve; }); + const json = vi.fn(() => body); + vi.stubGlobal("fetch", vi.fn(async () => ({ ok: true, json }))); + const pending = internal.request("/push"); + const rejected = expect(pending).rejects.toMatchObject({ code: "AUTH_SESSION_CHANGED" }); + while (!json.mock.calls.length) await Promise.resolve(); + localStorage.setItem("nowen-token", "new-account-token"); + finish({ results: [] }); + await rejected; + }); + it("blocks cross-account outbox requests during session reconfiguration", async () => { const { internal } = createEngine(); localStorage.setItem("nowen-token", "different-user"); diff --git a/frontend/src/lib/__tests__/mobileSyncRealtime.test.ts b/frontend/src/lib/__tests__/mobileSyncRealtime.test.ts new file mode 100644 index 000000000..0b7952b53 --- /dev/null +++ b/frontend/src/lib/__tests__/mobileSyncRealtime.test.ts @@ -0,0 +1,39 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import { installMobileSyncRealtime } from "../mobileSyncRealtime"; +import type { MobileSyncEngine } from "../mobileSyncEngine"; + +const mocks = vi.hoisted(() => ({ enabled: true, listeners: new Map void>(), connect: vi.fn() })); +vi.mock("../realtime", () => ({ realtime: { + connect: mocks.connect, + on: (event: string, callback: () => void) => { mocks.listeners.set(event, callback); return () => mocks.listeners.delete(event); }, +} })); +vi.mock("../mobileSyncStatus", () => ({ isMobileSyncEnabled: () => mocks.enabled })); +afterEach(() => { vi.useRealTimers(); mocks.listeners.clear(); mocks.enabled = true; mocks.connect.mockClear(); }); + +describe("native notification and fallback Pull scheduling", () => { + it("wakes the same engine on notice, reconnect, foreground and missing-notice poll", async () => { + vi.useFakeTimers(); + const requestSync = vi.fn(); + const dispose = installMobileSyncRealtime({ requestSync } as unknown as MobileSyncEngine); + mocks.listeners.get("sync.changed")!(); mocks.listeners.get("open")!(); + document.dispatchEvent(new Event("visibilitychange")); + await vi.advanceTimersByTimeAsync(30_000); + expect(requestSync).toHaveBeenCalledTimes(4); + expect(requestSync).toHaveBeenLastCalledWith(150); + expect(mocks.connect).toHaveBeenCalledOnce(); + dispose(); + await vi.advanceTimersByTimeAsync(60_000); + document.dispatchEvent(new Event("visibilitychange")); + expect(requestSync).toHaveBeenCalledTimes(4); + expect(mocks.listeners.size).toBe(0); + }); + it("does not synchronize or connect while the user has disabled sync", async () => { + vi.useFakeTimers(); mocks.enabled = false; + const requestSync = vi.fn(); + const dispose = installMobileSyncRealtime({ requestSync } as unknown as MobileSyncEngine); + mocks.listeners.get("sync.changed")!(); + await vi.advanceTimersByTimeAsync(30_000); + expect(requestSync).not.toHaveBeenCalled(); expect(mocks.connect).not.toHaveBeenCalled(); + dispose(); + }); +}); diff --git a/frontend/src/lib/__tests__/mobileSyncSwitch.test.ts b/frontend/src/lib/__tests__/mobileSyncSwitch.test.ts index 0c5838060..1d82862e7 100644 --- a/frontend/src/lib/__tests__/mobileSyncSwitch.test.ts +++ b/frontend/src/lib/__tests__/mobileSyncSwitch.test.ts @@ -106,7 +106,9 @@ beforeEach(async () => { if (url.endsWith("/scopes")) payload = { items: [{ scopeKey: "personal", workspaceId: null, workspaceName: null, role: null, canWrite: true, accessFingerprint: "test" }] }; else if (url.includes("/push?")) { const { mutations } = JSON.parse(String(init?.body)); - payload = { serverSequence: 1, results: mutations.map(({ mutationId }: { mutationId: string }) => ({ mutationId, status: "applied" })) }; + payload = { serverSequence: 1, results: mutations.map(({ mutationId, baseVersion, payload: note }: { mutationId: string; baseVersion?: number; payload?: {version?:number} }) => ({ + mutationId, status: "applied", version: baseVersion == null ? note?.version || 1 : baseVersion + 1, + })) }; } else if (url.includes("/changes?")) payload = { resetRequired: false, nextSequence: 1, items: [] }; else if (url.includes("/ack")) payload = {}; else if (url.includes("/blob/")) return { ok: init?.method !== "HEAD", status: init?.method === "HEAD" ? 404 : 200 }; diff --git a/frontend/src/lib/__tests__/voiceMemoLocalFirst.test.tsx b/frontend/src/lib/__tests__/voiceMemoLocalFirst.test.tsx index c1d300360..11f668a5a 100644 --- a/frontend/src/lib/__tests__/voiceMemoLocalFirst.test.tsx +++ b/frontend/src/lib/__tests__/voiceMemoLocalFirst.test.tsx @@ -1,6 +1,6 @@ import React, { act } from "react"; import { createRoot } from "react-dom/client"; -import { afterEach, describe, expect, it, vi } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import type { NativeDatabase } from "../nativeDatabase"; import type { NativeAttachmentStore } from "../nativeAttachmentStore"; @@ -46,6 +46,10 @@ function createDatabase() { } else if (/DELETE FROM sync_outbox WHERE mutationId/.test(sql)) { const index = outbox.findIndex((row) => row.mutationId === values[0]); if (index !== -1) outbox.splice(index, 1); + } else if (/UPDATE sync_outbox SET baseVersion/.test(sql)) { + const descendants = outbox.filter((row) => row.entityType === "note" && row.entityId === values[3]); + for (const row of descendants) row.baseVersion = values[0]; + return { changes: descendants.length }; } else if (/INSERT INTO native_runtime_meta/.test(sql) || /DELETE FROM native_runtime_meta/.test(sql)) { // Unsynced-note marker: the fake DB has no persistent native metadata table. } else throw new Error(`未覆盖的 SQL: ${sql}`); @@ -70,6 +74,7 @@ function createDatabase() { return { db, notes, outbox }; } let restore: (() => void) | undefined; +beforeEach(() => localStorage.clear()); afterEach(() => { restore?.(); restore = undefined; resetAttachmentAccessStateForTests(); vi.unstubAllGlobals(); }); describe("Android Local-first voice memo sync", () => { @@ -96,7 +101,9 @@ describe("Android Local-first voice memo sync", () => { const requests: Array<{ mutations: Array<{ mutationId: string; payload: { content: string } }> }> = []; vi.stubGlobal("fetch", vi.fn(async (_url: string, init: RequestInit) => { const body = JSON.parse(String(init.body)); requests.push(body); - return { ok: true, json: async () => ({ serverSequence: 2, results: body.mutations.map((mutation: { mutationId: string }) => ({ mutationId: mutation.mutationId, status: "applied" })) }) }; + return { ok: true, json: async () => ({ serverSequence: 2, results: body.mutations.map((mutation: { mutationId: string; baseVersion?: number; payload: {version:number} }) => ({ + mutationId: mutation.mutationId, status: "applied", version: mutation.baseVersion == null ? mutation.payload.version : mutation.baseVersion + 1, + })) }) }; })); const engine = new MobileSyncEngine({ db, attachments, serverUrl: "https://server.example", token: "test", userId: "user", profileId: "profile", deviceId: "device" }); await (engine as unknown as { push: (value: typeof scope) => Promise }).push(scope); diff --git a/frontend/src/lib/encryptedNotes/__tests__/nativeDatabaseInitialization.test.ts b/frontend/src/lib/encryptedNotes/__tests__/nativeDatabaseInitialization.test.ts index 865e4390d..98adcbcc0 100644 --- a/frontend/src/lib/encryptedNotes/__tests__/nativeDatabaseInitialization.test.ts +++ b/frontend/src/lib/encryptedNotes/__tests__/nativeDatabaseInitialization.test.ts @@ -39,3 +39,30 @@ it("a failed guard install rolls back without acknowledging a migrated schema", expect(statements.at(-1)).toBe("ROLLBACK"); expect(statements).not.toContain("PRAGMA user_version = 5"); expect(statements).not.toContain("COMMIT"); }); + +async function verifyConcurrentRead(rollback: boolean) { + const db = await openNativeDatabase("concurrent-sync-account"); + let release!: () => void, entered!: () => void; + const started = new Promise((resolve) => { entered = resolve; }); + const gate = new Promise((resolve) => { release = resolve; }); + let expired: typeof db | undefined; + const apply = db.transaction(async (tx) => { + expired = tx; + await tx.run("UPDATE notes SET title='fixture'"); + entered(); await gate; + if (rollback) throw new Error("apply failure"); + }); + const completion = apply.catch((error) => error.message); + await started; + // The reader belongs to another event/UI callback, not the transaction body. + let read: Promise; + try { read = db.query("SELECT id FROM notes"); } + finally { release(); } + expect(await completion).toBe(rollback ? "apply failure" : undefined); + expect(await read!).toEqual([]); + expect(mocks.execute.mock.calls.at(-1)?.[0]).toBe(rollback ? "ROLLBACK" : "COMMIT"); + expect(() => expired!.query("SELECT id FROM notes")).toThrow("事务视图已失效"); + await db.close(); +} +it("external UI reads wait for an in-progress Apply transaction to commit", () => verifyConcurrentRead(false)); +it("external UI reads wait for an in-progress Apply transaction to roll back", () => verifyConcurrentRead(true)); diff --git a/frontend/src/lib/mobileLocalFirstRuntime.ts b/frontend/src/lib/mobileLocalFirstRuntime.ts index 1603615f6..52edceaf2 100644 --- a/frontend/src/lib/mobileLocalFirstRuntime.ts +++ b/frontend/src/lib/mobileLocalFirstRuntime.ts @@ -12,9 +12,10 @@ import { } from "./localStore"; import { installMobileLocalFirstBridge } from "./mobileLocalFirstBridge"; import { migrateMobileLocalAccount } from "./mobileLocalAccountMigration"; -import { notifyMobileSyncStatusChanged, setMobileSyncEnabled } from "./mobileSyncStatus"; +import { notifyMobileSyncStatusChanged, setMobileSyncEnabled, MOBILE_SYNC_APPLIED_EVENT } from "./mobileSyncStatus"; import { createMobileSyncEngine, type MobileSyncEngine } from "./mobileSyncEngine"; import { createMobileServerEndpointRuntime } from "./mobileServerEndpointRuntime"; +import { installMobileSyncRealtime } from "./mobileSyncRealtime"; import { createNativeAttachmentStore } from "./nativeAttachmentStore"; import { openNativeDatabase, type NativeDatabase } from "./nativeDatabase"; import { createNativeLocalRepository } from "./nativeLocalRepository"; @@ -546,6 +547,9 @@ async function configureRuntime(): Promise { beforeSync: () => endpoints?.beforeSync() || Promise.resolve(), onNetworkUnavailable: () => endpoints?.recover(), onAttachmentReady:(id)=>{void repository?.refreshAttachmentUrl(id);}, + onApplied: (detail) => { + if (active?.engine === engine) window.dispatchEvent(new CustomEvent(MOBILE_SYNC_APPLIED_EVENT, { detail })); + }, }); repository = createNativeLocalRepository({ db,attachments:attachmentStore,accountId:ids.accountId,userId, @@ -565,6 +569,7 @@ async function configureRuntime(): Promise { active = {identity,db,engine,restoreBridge,removeListeners}; const profile = (await db.query<{ enabled: number }>("SELECT enabled FROM sync_profiles WHERE id=?",[ids.profileId]))[0]; setMobileSyncEnabled(profile?.enabled === 1); + removeListeners.push(installMobileSyncRealtime(engine)); if (profile?.enabled === 1) engine.start(); } catch (error) { for (const remove of removeListeners) await Promise.resolve(remove()).catch(() => undefined); diff --git a/frontend/src/lib/mobileNoteSyncReceipt.ts b/frontend/src/lib/mobileNoteSyncReceipt.ts index 5f70bf9f3..82821af2b 100644 --- a/frontend/src/lib/mobileNoteSyncReceipt.ts +++ b/frontend/src/lib/mobileNoteSyncReceipt.ts @@ -30,9 +30,10 @@ export async function readNativeNoteSyncReceipt( ); try { const receipt = rows[0] && JSON.parse(rows[0].value) as { mutationId?: unknown; serverVersion?: unknown }; - if (typeof receipt?.mutationId === "string" && receipt.mutationId) { - const version = typeof receipt.serverVersion === "number" && Number.isSafeInteger(receipt.serverVersion) - ? receipt.serverVersion : null; + if (typeof receipt?.mutationId === "string" && receipt.mutationId + && typeof receipt.serverVersion === "number" && Number.isSafeInteger(receipt.serverVersion) + && receipt.serverVersion > 0) { + const version = receipt.serverVersion; return { phase: "confirmed", localRevision: null, acknowledgedRevision: version }; } } catch { /* Corrupt metadata cannot prove an ACK. */ } diff --git a/frontend/src/lib/mobileSyncEngine.ts b/frontend/src/lib/mobileSyncEngine.ts index 063b6964f..af7b847ac 100644 --- a/frontend/src/lib/mobileSyncEngine.ts +++ b/frontend/src/lib/mobileSyncEngine.ts @@ -5,7 +5,7 @@ import { getResolvedApiBaseUrl } from "./serverUrl"; import { fetchWithAuthRefresh, getAccessToken } from "./authSession"; import { SERVER_ENDPOINT_CHANGED_EVENT } from "./serverEndpointState"; import { forgetUnsentLocalNotes } from "./nativeLocalNoteOrigin"; -import { notifyMobileSyncStatusChanged } from "./mobileSyncStatus"; +import { notifyMobileSyncStatusChanged, type MobileSyncAppliedDetail } from "./mobileSyncStatus"; import { nativeReceiptKey } from "./mobileNoteSyncReceipt"; import { validateEncryptedNoteWrite } from "./encryptedNotes/noteDocument"; @@ -35,6 +35,7 @@ interface MobileSyncOptions { onAttachmentReady?: (attachmentId:string) => void; beforeSync?: () => Promise; onNetworkUnavailable?: () => void; + onApplied?: (detail: MobileSyncAppliedDetail) => void; } interface SnapshotEntry { @@ -180,10 +181,14 @@ export class MobileSyncEngine { ))[0]; if (state?.accessStatus === "access_revoked") continue; try { + // Replay durable mutation IDs before reconciling a snapshot. After a + // lost ACK, the server already has our create; bootstrap alone cannot + // distinguish that state from a competing creation on another device. + await this.push(scope); + if (this.stopped) return; if (!state || state.lastSequence === 0 || state.accessStatus === "replan_required") { await this.bootstrap(scope); } - await this.push(scope); if (this.stopped) return; await this.pull(scope); if (this.stopped) return; @@ -222,7 +227,7 @@ export class MobileSyncEngine { } finally { this.running = false; this.syncAbort = null; - notifyMobileSyncStatusChanged(); + if (!this.stopped) notifyMobileSyncStatusChanged(); if (this.rerun) { this.rerun = false; this.requestSync(0); } } } @@ -232,15 +237,17 @@ export class MobileSyncEngine { resolution:"keep-local"|"keep-remote"|"manual", mergedPayload?:Record, ):Promise{ + this.syncAbort?.abort(); const row=(await this.options.db.query<{ id:string;scopeKey:string;entityType:RemoteEntityType;entityId:string; localVersion:number|null;remoteVersion:number|null;localPayload:string|null;remotePayload:string|null;status:string; - }>("SELECT * FROM sync_conflicts WHERE id=?",[conflictId]))[0]; + }>("SELECT * FROM sync_conflicts WHERE id=? AND profileId=?",[conflictId,this.options.profileId]))[0]; if(!row||row.status==="resolved")return; if(!isCoreEntityType(row.entityType))throw new Error("该冲突类型暂不支持在移动端合并"); const parse=(value:string|null)=>value?JSON.parse(value) as Record:null; - const payload=resolution==="manual"?mergedPayload:parse(resolution==="keep-remote"?row.remotePayload:row.localPayload); - if(!payload)throw new Error("缺少可用的冲突内容"); + const selected=resolution==="manual"?mergedPayload:parse(resolution==="keep-remote"?row.remotePayload:row.localPayload); + if(!selected)throw new Error("缺少可用的冲突内容"); + if(resolution!=="keep-remote"&&parse(row.remotePayload)?.__delete)throw new Error("远端笔记已永久删除,请先另存为新笔记保留本机内容"); const workspace=workspaceId(row.scopeKey); const stored=(await this.options.db.query<{workspaceName:string;role:string;canWrite:number;accessFingerprint:string}>( "SELECT workspaceName,role,canWrite,accessFingerprint FROM sync_workspace_scopes WHERE profileId=? AND scopeKey=?", @@ -249,6 +256,13 @@ export class MobileSyncEngine { const scope:ScopeDescriptor={scopeKey:row.scopeKey,workspaceId:workspace,workspaceName:stored?.workspaceName||null, role:stored?.role||null,canWrite:workspace?stored?.canWrite===1:true,accessFingerprint:stored?.accessFingerprint||""}; await this.options.db.transaction(async(tx)=>{ + // The editor can keep saving while the conflict is awaiting a choice. + // Keep-local must preserve that latest durable edit, not its old snapshot. + const chosen = resolution === "keep-local" + ? await this.readEntity(tx, row.scopeKey, row.entityType, row.entityId) || selected : selected; + const payload = row.entityType === "note" && resolution !== "keep-remote" + && typeof row.remoteVersion === "number" && Number.isSafeInteger(row.remoteVersion) && row.remoteVersion > 0 + ? { ...chosen, version: row.remoteVersion + 1 } : chosen; await tx.run(`DELETE FROM sync_outbox WHERE profileId=? AND scopeKey=? AND entityType=? AND entityId=? AND status IN ('pending','inflight','failed')`, [this.options.profileId,row.scopeKey,row.entityType,row.entityId]); @@ -267,8 +281,8 @@ export class MobileSyncEngine { } async forkLocalConflict(conflictId:string,side:"local"|"remote"):Promise{ - const row=(await this.options.db.query<{scopeKey:string;entityType:RemoteEntityType;localPayload:string|null;remotePayload:string|null}>( - "SELECT scopeKey,entityType,localPayload,remotePayload FROM sync_conflicts WHERE id=?",[conflictId], + const row=(await this.options.db.query<{scopeKey:string;entityId:string;entityType:RemoteEntityType;localPayload:string|null;remotePayload:string|null}>( + "SELECT scopeKey,entityId,entityType,localPayload,remotePayload FROM sync_conflicts WHERE id=? AND profileId=?",[conflictId,this.options.profileId], ))[0]; if(!row||row.entityType!=="note")throw new Error("只有笔记冲突支持另存为新笔记"); const payload=JSON.parse(side==="local"?row.localPayload||"null":row.remotePayload||"null") as Record|null; @@ -280,8 +294,9 @@ export class MobileSyncEngine { ))[0]; const scope:ScopeDescriptor={scopeKey:row.scopeKey,workspaceId:workspace,workspaceName:stored?.workspaceName||null, role:stored?.role||null,canWrite:workspace?stored?.canWrite===1:true,accessFingerprint:stored?.accessFingerprint||""}; - const next={...payload,id,workspaceId:workspace,title:`${String(payload.title||"无标题笔记")}(${side==="local"?"本机":"服务器"}版本)`,version:1}; await this.options.db.transaction(async(tx)=>{ + const chosen = side === "local" ? await this.readEntity(tx, row.scopeKey, "note", row.entityId) || payload : payload; + const next={...chosen,id,workspaceId:workspace,title:`${String(chosen.title||"无标题笔记")}(${side==="local"?"本机":"服务器"}版本)`,version:1}; await this.applyEntity(tx,scope,{entityType:"note",entityId:id,payload:next}); await tx.run(`INSERT INTO sync_outbox ( id,mutationId,profileId,deviceId,scopeKey,entityType,entityId,operation,payload,status,retryCount,createdAt @@ -314,6 +329,7 @@ export class MobileSyncEngine { init.signal?.addEventListener("abort", cancel, { once: true }); syncSignal?.addEventListener("abort", cancel, { once: true }); window.addEventListener(SERVER_ENDPOINT_CHANGED_EVENT, endpointChanged); + if (syncSignal) window.addEventListener("nowen:token-changed", cancel); if (init.signal?.aborted) cancel(); try { // Capacitor 的原生 POST/PUT fetch 不消费 AbortSignal;仍须释放同步循环, @@ -329,7 +345,10 @@ export class MobileSyncEngine { ...(init.body ? { "Content-Type": "application/json" } : {}), ...(init.headers || {}), }, - }, apiBaseUrl, fetch)]); + }, apiBaseUrl, (input, requestInit) => { + controller.signal.throwIfAborted(); + return fetch(input, requestInit); + })]); } catch (cause) { if (!syncSignal?.aborted) this.options.onNetworkUnavailable?.(); throw Object.assign(new Error("网络不可用", { cause }), { code: "NETWORK_UNAVAILABLE" }); @@ -339,19 +358,24 @@ export class MobileSyncEngine { init.signal?.removeEventListener("abort", cancel); syncSignal?.removeEventListener("abort", cancel); window.removeEventListener(SERVER_ENDPOINT_CHANGED_EVENT, endpointChanged); + if (syncSignal) window.removeEventListener("nowen:token-changed", cancel); } } private async request(path: string, init: RequestInit = {}): Promise { + const signal = this.syncAbort?.signal; const response = await this.fetchResponse(path, init); + signal?.throwIfAborted(); + const responseToken = getAccessToken(); if (!response.ok) { const payload = await response.json().catch(() => ({})) as { code?: string; error?: string }; // A reverse proxy may use a different JSON error code for HTTP 401. const code = response.status === 401 ? "AUTH_EXPIRED" : (payload.code || "SERVER_ERROR"); throw Object.assign(new Error(payload.error || code), { code, status: response.status }); } + let payload: T; try { - return await response.json() as T; + payload = await response.json() as T; } catch { // NodeBao/reverse proxies sometimes return their HTML login/home page // with HTTP 200 for an unrecognized /sync/v2 route. This is not a @@ -361,6 +385,11 @@ export class MobileSyncEngine { "同步接口返回了非 JSON 响应,请检查节点小宝的 API 代理路径与登录状态", ); } + signal?.throwIfAborted(); + if (getAccessToken() !== responseToken) { + throw syncFailure("AUTH_SESSION_CHANGED", "登录状态已更新,旧同步响应已隔离"); + } + return payload; } private async fetchScopes(): Promise { @@ -453,39 +482,51 @@ export class MobileSyncEngine { const stale=await this.options.db.query<{id:string}>(`SELECT id FROM attachments WHERE scopeKey=? ${attachmentIds.length?`AND id NOT IN (${attachmentIds.map(()=>"?").join(",")})`:""}`,[scope.scopeKey,...attachmentIds]); await this.options.db.transaction(async(tx)=>{ - const remove=async(table:string,type:EntityType)=>{ + const remove=async(table:string,type:EntityType,extra="")=>{ const ids=[...(seen.get(type)||[])]; await tx.run(`DELETE FROM ${table} WHERE scopeKey=? ${ids.length?`AND id NOT IN (${ids.map(()=>"?").join(",")})`:""} AND NOT EXISTS (SELECT 1 FROM sync_outbox o WHERE o.profileId=? AND o.scopeKey=? - AND o.entityType=? AND o.entityId=${table}.id AND o.status IN ('pending','inflight','failed'))`, - [scope.scopeKey,...ids,this.options.profileId,scope.scopeKey,type]); + AND o.entityType=? AND o.entityId=${table}.id AND o.status IN ('pending','inflight','failed')) + AND NOT EXISTS (SELECT 1 FROM sync_conflicts c WHERE c.profileId=? AND c.scopeKey=? + AND c.entityType=? AND c.entityId=${table}.id AND c.status='unresolved') ${extra}`, + [scope.scopeKey,...ids,this.options.profileId,scope.scopeKey,type,this.options.profileId,scope.scopeKey,type]); }; const removeComposite=async(table:string,type:EntityType,idSql:string)=>{ const ids=[...(seen.get(type)||[])]; await tx.run(`DELETE FROM ${table} WHERE scopeKey=? ${ids.length?`AND (${idSql}) NOT IN (${ids.map(()=>"?").join(",")})`:""} AND NOT EXISTS (SELECT 1 FROM sync_outbox o WHERE o.profileId=? AND o.scopeKey=? - AND o.entityType=? AND o.entityId=(${idSql}) AND o.status IN ('pending','inflight','failed'))`, - [scope.scopeKey,...ids,this.options.profileId,scope.scopeKey,type]); + AND o.entityType=? AND o.entityId=(${idSql}) AND o.status IN ('pending','inflight','failed')) + AND NOT EXISTS (SELECT 1 FROM sync_conflicts c WHERE c.profileId=? AND c.scopeKey=? + AND c.entityType=? AND c.entityId=(${idSql}) AND c.status='unresolved')`, + [scope.scopeKey,...ids,this.options.profileId,scope.scopeKey,type,this.options.profileId,scope.scopeKey,type]); }; - await remove("attachments","attachment"); + await remove("attachments","attachment","AND NOT (available=1 AND transferStatus IN ('local','pending_upload','uploading','failed'))"); await remove("diaries","diary"); await remove("mindmaps","mindmap"); const reminderIds=[...(seen.get("task_reminder")||[])]; await tx.run(`DELETE FROM task_reminders WHERE taskId IN (SELECT id FROM tasks WHERE scopeKey=?) ${reminderIds.length?`AND id NOT IN (${reminderIds.map(()=>"?").join(",")})`:""} AND NOT EXISTS (SELECT 1 FROM sync_outbox o WHERE o.profileId=? AND o.scopeKey=? - AND o.entityType='task_reminder' AND o.entityId=task_reminders.id AND o.status IN ('pending','inflight','failed'))`, - [scope.scopeKey,...reminderIds,this.options.profileId,scope.scopeKey]); - await remove("tasks","task"); + AND o.entityType='task_reminder' AND o.entityId=task_reminders.id AND o.status IN ('pending','inflight','failed')) + AND NOT EXISTS (SELECT 1 FROM sync_conflicts c WHERE c.profileId=? AND c.scopeKey=? + AND c.entityType='task_reminder' AND c.entityId=task_reminders.id AND c.status='unresolved')`, + [scope.scopeKey,...reminderIds,this.options.profileId,scope.scopeKey,this.options.profileId,scope.scopeKey]); + await remove("tasks","task","AND NOT EXISTS (SELECT 1 FROM task_reminders r WHERE r.taskId=tasks.id)"); await removeComposite("favorites","favorite","favorites.userId || ':' || favorites.noteId"); await removeComposite("note_tags","note_tag","note_tags.noteId || ':' || note_tags.tagId"); - await remove("notes","note"); - await remove("notebooks","notebook"); - await remove("tags","tag"); + await remove("notes","note",`AND NOT EXISTS (SELECT 1 FROM attachments a WHERE a.scopeKey=notes.scopeKey AND a.noteId=notes.id) + AND NOT EXISTS (SELECT 1 FROM note_tags nt WHERE nt.scopeKey=notes.scopeKey AND nt.noteId=notes.id) + AND NOT EXISTS (SELECT 1 FROM favorites f WHERE f.scopeKey=notes.scopeKey AND f.noteId=notes.id)`); + await remove("notebooks","notebook",`AND NOT EXISTS (SELECT 1 FROM notes n WHERE n.scopeKey=notebooks.scopeKey AND n.notebookId=notebooks.id) + AND NOT EXISTS (SELECT 1 FROM notebooks child WHERE child.scopeKey=notebooks.scopeKey AND child.parentId=notebooks.id)`); + await remove("tags","tag","AND NOT EXISTS (SELECT 1 FROM note_tags nt WHERE nt.scopeKey=tags.scopeKey AND nt.tagId=tags.id)"); }); - await Promise.all(stale.map(({id})=>this.options.attachments.remove(id).catch(()=>undefined))); + for (const { id } of stale) { + const retained = await this.options.db.query("SELECT id FROM attachments WHERE scopeKey=? AND id=?", [scope.scopeKey,id]); + if (!retained.length) await this.options.attachments.remove(id); + } } private async push(scope: ScopeDescriptor): Promise { @@ -493,7 +534,11 @@ export class MobileSyncEngine { mutationId:string;entityType:EntityType;entityId:string;operation:"upsert"|"delete"; baseVersion:number|null;payload:string|null; }>(`SELECT mutationId,entityType,entityId,operation,baseVersion,payload FROM sync_outbox - WHERE profileId=? AND scopeKey=? AND status IN ('pending','failed') ORDER BY createdAt LIMIT 100`, + WHERE profileId=? AND scopeKey=? AND status IN ('pending','failed') + AND NOT EXISTS (SELECT 1 FROM sync_conflicts c WHERE c.profileId=sync_outbox.profileId + AND c.scopeKey=sync_outbox.scopeKey AND c.entityType=sync_outbox.entityType + AND c.entityId=sync_outbox.entityId AND c.status='unresolved') + ORDER BY createdAt,rowid LIMIT 100`, [this.options.profileId,scope.scopeKey]); if (!rows.length || !scope.canWrite && scope.workspaceId) return; for(const row of rows){ @@ -502,6 +547,18 @@ export class MobileSyncEngine { } const mutations = rows.map((row) => ({ ...row, baseVersion:row.baseVersion??undefined,payload:row.payload?JSON.parse(row.payload):undefined })); + // Acknowledgements rebase all unsent descendants to the confirmed version. + // The server applies this batch in order, so later note mutations must then + // use the revision produced by the preceding mutation of the same note. + const noteVersions = new Map(); + for (const mutation of mutations) { + if (mutation.entityType !== "note") continue; + if (noteVersions.has(mutation.entityId)) mutation.baseVersion = noteVersions.get(mutation.entityId); + if (mutation.operation === "upsert") { + noteVersions.set(mutation.entityId, mutation.baseVersion === undefined + ? Math.max(1, Number(mutation.payload?.version) || 1) : mutation.baseVersion + 1); + } else noteVersions.delete(mutation.entityId); + } for (const mutation of mutations) if (mutation.entityType === "note" && mutation.operation === "upsert") { validateEncryptedNoteWrite(mutation.payload || {}); const current = await this.readEntity(this.options.db, scope.scopeKey, "note", mutation.entityId); @@ -516,7 +573,13 @@ export class MobileSyncEngine { ); for (const result of response.results) { const source = mutations.find((mutation) => mutation.mutationId === result.mutationId); - if (source?.entityType === "note" && result.code === "VERSION_CONFLICT" && result.serverPayload) { + if (!source) throw syncFailure("SERVER_ERROR", "同步 ACK 不属于当前发送批次"); + if (source.entityType === "note" && source.operation === "upsert" + && (result.status === "applied" || result.status === "duplicate") + && !(typeof result.version === "number" && Number.isSafeInteger(result.version) && result.version > 0)) { + throw syncFailure("SERVER_ERROR", "笔记同步 ACK 缺少有效版本,已保留待发送修改"); + } + if (source.entityType === "note" && result.code === "VERSION_CONFLICT" && result.serverPayload && !result.serverPayload.__delete) { validateEncryptedNoteWrite(result.serverPayload); validateEncryptedNoteWrite(result.serverPayload, source.payload); } @@ -527,13 +590,17 @@ export class MobileSyncEngine { if (result.status === "applied" || result.status === "duplicate") { await tx.run("DELETE FROM sync_outbox WHERE mutationId=?",[result.mutationId]); if (source?.entityType === "note" && source.operation === "upsert") { + // Edits committed during Push descend from this server revision. + // Keep their content and calibrate the next batch's base atomically. + const remaining = await tx.run(`UPDATE sync_outbox SET baseVersion=? + WHERE profileId=? AND scopeKey=? AND entityType='note' AND entityId=? + AND status IN ('pending','failed')`, [result.version, this.options.profileId, scope.scopeKey, source.entityId]); + if (!remaining.changes) await tx.run("UPDATE notes SET version=? WHERE scopeKey=? AND id=?", [result.version, scope.scopeKey, source.entityId]); // Persist the mutation-specific ACK atomically with the Outbox removal. await tx.run(`INSERT INTO native_runtime_meta (key,value,updatedAt) VALUES (?,?,?) ON CONFLICT(key) DO UPDATE SET value=excluded.value,updatedAt=excluded.updatedAt`, [ nativeReceiptKey(this.options.profileId, source.entityId), - JSON.stringify({ mutationId: result.mutationId, - serverVersion: typeof result.version === "number" && Number.isSafeInteger(result.version) && result.version > 0 - ? result.version : null }), + JSON.stringify({ mutationId: result.mutationId, serverVersion: result.version }), now(), ]); } @@ -541,12 +608,13 @@ export class MobileSyncEngine { await this.options.attachments.remove(source.entityId).catch(() => undefined); } } else if (result.code === "VERSION_CONFLICT" && source) { + const local = await this.readEntity(tx, scope.scopeKey, source.entityType, source.entityId); await tx.run(`INSERT INTO sync_conflicts ( id,profileId,scopeKey,entityType,entityId,localVersion,remoteVersion, localPayload,remotePayload,status,createdAt ) VALUES (?,?,?,?,?,?,?,?,?,'unresolved',?)`, [ newLocalId(),this.options.profileId,scope.scopeKey,source.entityType,source.entityId, - source.baseVersion,result.serverVersion,source.payload, + source.baseVersion,result.serverVersion ?? null,local ? JSON.stringify(local) : source.payload, JSON.stringify(result.serverPayload||{version:result.serverVersion}),now(), ]); await tx.run("DELETE FROM sync_outbox WHERE mutationId=?",[result.mutationId]); @@ -570,9 +638,12 @@ export class MobileSyncEngine { ); assertMobileSyncItems(changes.items,"changes"); if (changes.resetRequired) { await this.bootstrap(scope); return; } - const deletes: SnapshotEntry[] = changes.items.filter((item)=>item.operation==="delete") + // Snapshot contains current entities, so intermediate upserts of an entity + // already deleted in this feed page must not be requested. + const latest = [...new Map(changes.items.map((item) => [`${item.entityType}\0${item.entityId}`, item])).values()]; + const deletes: SnapshotEntry[] = latest.filter((item)=>item.operation==="delete") .map((item)=>({entityType:item.entityType,entityId:item.entityId,payload:{__delete:true}})); - const wanted = new Set(changes.items.filter((item)=>item.operation!=="delete").map((item)=>`${item.entityType}\0${item.entityId}`)); + const wanted = new Set(latest.filter((item)=>item.operation!=="delete").map((item)=>`${item.entityType}\0${item.entityId}`)); const entries: SnapshotEntry[] = []; if (wanted.size) { let cursor: string|null=null; @@ -607,6 +678,7 @@ export class MobileSyncEngine { for (const entry of entries) if (entry.entityType === "note" && !entry.payload.__delete) validateEncryptedNoteWrite(entry.payload); const notebookParents=entries.flatMap((entry)=>entry.entityType==="notebook" && typeof entry.payload.parentId==="string" ? [[entry.entityId,entry.payload.parentId] as const] : []); + const applied: SnapshotEntry[] = []; await this.options.db.transaction(async (tx) => { // 即使远端内容进入冲突中心而未覆盖本地,它的存在也已得到确认。 await forgetUnsentLocalNotes(tx, scope.scopeKey, entries.filter((entry) => entry.entityType === "note").map((entry) => entry.entityId)); @@ -618,17 +690,17 @@ export class MobileSyncEngine { const conflict=(await tx.query<{id:string}>(`SELECT id FROM sync_conflicts WHERE profileId=? AND scopeKey=? AND entityType=? AND entityId=? AND status='unresolved' LIMIT 1`, [this.options.profileId,scope.scopeKey,entry.entityType,entry.entityId]))[0]; - if(conflict&&!entry.payload.__delete){ + if(conflict){ await tx.run(`UPDATE sync_conflicts SET remotePayload=?,remoteVersion=COALESCE(?,remoteVersion) WHERE id=?`,[JSON.stringify(entry.payload),Number(entry.payload.version)||null,conflict.id]); continue; } const pending=(await tx.query<{payload:string|null;baseVersion:number|null}>(`SELECT payload,baseVersion FROM sync_outbox WHERE - profileId=? AND scopeKey=? AND entityType=? AND entityId=? AND status IN ('pending','failed','inflight') LIMIT 1`, + profileId=? AND scopeKey=? AND entityType=? AND entityId=? AND status IN ('pending','failed','inflight') ORDER BY createdAt DESC,rowid DESC LIMIT 1`, [this.options.profileId,scope.scopeKey,entry.entityType,entry.entityId]))[0]; - if(pending&&!entry.payload.__delete){ + if(pending){ const local=await this.readEntity(tx,scope.scopeKey,entry.entityType,entry.entityId); - if (!bootstrap && entry.entityType === "note" && local) { + if (entry.entityType === "note" && local && !entry.payload.__delete) { // A previous mutation from THIS profile was already acknowledged, // and the pull sees that exact server version while a newer local // edit waits in Outbox. This is our own echo, not a true conflict. @@ -639,22 +711,30 @@ export class MobileSyncEngine { ))[0]; if (isAckedOwnNoteEcho(entry.payload, local, pending, receipt?.value)) continue; } - if(!bootstrap||stable(local)!==stable(entry.payload)){ + if(entry.payload.__delete||!bootstrap||stable(local)!==stable(entry.payload)){ await tx.run(`INSERT INTO sync_conflicts (id,profileId,scopeKey,entityType,entityId, - localPayload,remotePayload,status,createdAt) VALUES (?,?,?,?,?,?,?,'unresolved',?)`,[ + localPayload,remotePayload,localVersion,remoteVersion,status,createdAt) VALUES (?,?,?,?,?,?,?,?,?,'unresolved',?)`,[ newLocalId(),this.options.profileId,scope.scopeKey,entry.entityType,entry.entityId, - pending.payload||JSON.stringify(local),JSON.stringify(entry.payload),now(), + local ? JSON.stringify(local) : pending.payload,JSON.stringify(entry.payload), + typeof local?.version === "number" ? local.version : null, + typeof entry.payload.version === "number" ? entry.payload.version : null,now(), ]); continue; } } await this.applyEntity(tx,scope,entry); + applied.push(entry); } for(const [id,parentId] of notebookParents){ const parent=(await tx.query<{id:string}>("SELECT id FROM notebooks WHERE scopeKey=? AND id=?",[scope.scopeKey,parentId]))[0]; if(parent)await tx.run("UPDATE notebooks SET parentId=? WHERE scopeKey=? AND id=?",[parentId,scope.scopeKey,id]); } }); + if (applied.length && !this.stopped) this.options.onApplied?.({ + scopeKey: scope.scopeKey, + noteIds: applied.filter((entry) => entry.entityType === "note" && !entry.payload.__delete).map((entry) => entry.entityId), + deletedNoteIds: applied.filter((entry) => entry.entityType === "note" && entry.payload.__delete).map((entry) => entry.entityId), + }); } private async readEntity(db:NativeDatabase,scopeKey:string,type:EntityType,id:string):Promise|null>{ diff --git a/frontend/src/lib/mobileSyncRealtime.ts b/frontend/src/lib/mobileSyncRealtime.ts new file mode 100644 index 000000000..46108c5dc --- /dev/null +++ b/frontend/src/lib/mobileSyncRealtime.ts @@ -0,0 +1,19 @@ +import { realtime } from "./realtime"; +import { isMobileSyncEnabled } from "./mobileSyncStatus"; +import type { MobileSyncEngine } from "./mobileSyncEngine"; + +/** One notification transport; durable data always comes from Sync V2 Pull. */ +export function installMobileSyncRealtime(engine: MobileSyncEngine): () => void { + const wake = () => { if (isMobileSyncEnabled()) engine.requestSync(150); }; + const removeNotice = realtime.on("sync.changed", wake); + const removeOpen = realtime.on("open", wake); + const onVisible = () => { if (document.visibilityState === "visible") wake(); }; + document.addEventListener("visibilitychange", onVisible); + const poll = setInterval(onVisible, 30_000); + if (isMobileSyncEnabled()) realtime.connect(); + return () => { + clearInterval(poll); + removeNotice(); removeOpen(); + document.removeEventListener("visibilitychange", onVisible); + }; +} diff --git a/frontend/src/lib/mobileSyncStatus.ts b/frontend/src/lib/mobileSyncStatus.ts index 91eec5b16..b8b17898f 100644 --- a/frontend/src/lib/mobileSyncStatus.ts +++ b/frontend/src/lib/mobileSyncStatus.ts @@ -1,5 +1,11 @@ export const MOBILE_SYNC_STATUS_CHANGED_EVENT = "nowen:mobile-sync-status-changed"; export const MOBILE_SYNC_SETTINGS_CHANGED_EVENT = "nowen:mobile-sync-settings-changed"; +export const MOBILE_SYNC_APPLIED_EVENT = "nowen:mobile-sync-applied"; +export interface MobileSyncAppliedDetail { + scopeKey: string; + noteIds: string[]; + deletedNoteIds: string[]; +} let syncEnabled = false; export function isMobileSyncEnabled(): boolean { return syncEnabled; } diff --git a/frontend/src/lib/nativeDatabase.ts b/frontend/src/lib/nativeDatabase.ts index 4e1a54288..f99c35c03 100644 --- a/frontend/src/lib/nativeDatabase.ts +++ b/frontend/src/lib/nativeDatabase.ts @@ -7,6 +7,8 @@ import { } from "@capacitor-community/sqlite"; export interface NativeDatabase { + // Outer calls are serialized behind active transactions. A transaction body + // must use its tx view; awaiting the outer queue inside it would deadlock. run( sql: string, values?: unknown[], @@ -665,7 +667,6 @@ async function hashAccountId(accountId: string): Promise { class NativeDatabaseImpl implements NativeDatabase { private queue: Promise = Promise.resolve(); - private transactionActive = false; private closing = false; private closed = false; private closePromise?: Promise; @@ -758,7 +759,6 @@ class NativeDatabaseImpl implements NativeDatabase { this.assertAcceptingWork(); return this.enqueue(async () => { const transactionView = this.createTransactionView(); - this.transactionActive = true; try { return await this.withImmediateTransaction(async () => { try { @@ -769,13 +769,11 @@ class NativeDatabaseImpl implements NativeDatabase { }); } finally { transactionView.deactivate(); - this.transactionActive = false; } }); } close(): Promise { - this.assertOuterHandleAvailable(); if (this.closed) return Promise.resolve(); if (this.closePromise) return this.closePromise; @@ -806,19 +804,10 @@ class NativeDatabaseImpl implements NativeDatabase { } private assertAcceptingWork(): void { - this.assertOuterHandleAvailable(); if (this.closed) throw new Error("原生数据库连接已关闭"); if (this.closing) throw new Error("原生数据库连接正在关闭"); } - private assertOuterHandleAvailable(): void { - if (this.transactionActive) { - throw new Error( - "原生数据库事务执行期间不能使用外层句柄;请改用事务回调传入的 tx", - ); - } - } - private enqueue(work: () => Promise): Promise { const scheduled = this.queue.then(work, work); this.queue = scheduled.then( diff --git a/frontend/src/lib/nativeLocalRepository.ts b/frontend/src/lib/nativeLocalRepository.ts index a526824b8..0431f3f8c 100644 --- a/frontend/src/lib/nativeLocalRepository.ts +++ b/frontend/src/lib/nativeLocalRepository.ts @@ -244,9 +244,9 @@ export class NativeLocalRepository implements LocalRepository { return { scopeKey, workspaceId: workspaceIdFromScope(scopeKey) }; } - private async assertWritable(scopeKey:string):Promise { + private async assertWritable(scopeKey:string, db:NativeDatabase=this.db):Promise { if(scopeKey === "personal") return; - const scope=(await this.db.query<{canWrite:number;accessStatus:string}>( + const scope=(await db.query<{canWrite:number;accessStatus:string}>( "SELECT canWrite,accessStatus FROM sync_workspace_scopes WHERE scopeKey=? LIMIT 1",[scopeKey], ))[0]; if(!scope || scope.canWrite !== 1 || scope.accessStatus !== "active") { @@ -378,13 +378,13 @@ export class NativeLocalRepository implements LocalRepository { `, values); } - private async getNote(id: string): Promise { + private async getNote(id: string, db: NativeDatabase = this.db): Promise { const currentScopeKey = this.scope().scopeKey; - const rows = await this.db.query(`SELECT * FROM notes WHERE id=? + const rows = await db.query(`SELECT * FROM notes WHERE id=? ORDER BY CASE WHEN scopeKey=? THEN 0 ELSE 1 END LIMIT 1`, [id,currentScopeKey]); if (!rows[0]) return null; const scopeKey = this.scopeFromWorkspace(rows[0].workspaceId).scopeKey; - const tags = await this.db.query(` + const tags = await db.query(` SELECT t.* FROM tags t JOIN note_tags nt ON nt.scopeKey=t.scopeKey AND nt.tagId=t.id WHERE nt.scopeKey=? AND nt.noteId=? ORDER BY t.name @@ -416,7 +416,7 @@ export class NativeLocalRepository implements LocalRepository { ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)`, Object.values(row)); await tx.run(`INSERT INTO native_runtime_meta (key,value,updatedAt) VALUES (?,'1',?) ON CONFLICT(key) DO UPDATE SET value=excluded.value,updatedAt=excluded.updatedAt`, [unsentLocalNoteKey(scope.scopeKey, input.id), savedAt]); - await this.enqueue(tx, "note", input.id, "upsert", row); + await this.enqueue(tx, "note", input.id, "upsert", row, undefined, scope.scopeKey); }); // The receipt is invalidated only after the SQLite Outbox transaction commits. notifyMobileSyncStatusChanged(); @@ -424,16 +424,16 @@ export class NativeLocalRepository implements LocalRepository { } private async updateNote(id: string, patch: Partial): Promise { - const current = await this.getNote(id); - if (!current) throw new Error("笔记不存在"); - validateEncryptedNoteWrite(patch as Record, current as unknown as Record); - if (isEncryptedNoteFormat(current.contentFormat) && patch.version !== undefined && patch.version !== current.version) throw Object.assign(new Error("笔记版本已改变"), { status: 409, code: "VERSION_CONFLICT" }); - const encrypted = isEncryptedNoteFormat(current.contentFormat); - if (encrypted && patch.content !== undefined && !Number.isSafeInteger(patch.version)) throw Object.assign(new Error("加密正文保存需要明确版本"), { status: 409, code: "VERSION_CONFLICT" }); - const next = { ...current, ...patch, id, updatedAt: now(), version: current.version + 1 }; - const scope = this.scopeFromWorkspace(current.workspaceId); - await this.assertWritable(scope.scopeKey); - await this.db.transaction(async (tx) => { + const result = await this.db.transaction(async (tx) => { + const current = await this.getNote(id, tx); + if (!current) throw new Error("笔记不存在"); + validateEncryptedNoteWrite(patch as Record, current as unknown as Record); + if (isEncryptedNoteFormat(current.contentFormat) && patch.version !== undefined && patch.version !== current.version) throw Object.assign(new Error("笔记版本已改变"), { status: 409, code: "VERSION_CONFLICT" }); + const encrypted = isEncryptedNoteFormat(current.contentFormat); + if (encrypted && patch.content !== undefined && !Number.isSafeInteger(patch.version)) throw Object.assign(new Error("加密正文保存需要明确版本"), { status: 409, code: "VERSION_CONFLICT" }); + const next = { ...current, ...patch, id, updatedAt: now(), version: current.version + 1 }; + const scope = this.scopeFromWorkspace(current.workspaceId); + await this.assertWritable(scope.scopeKey, tx); const changed = await tx.run(`UPDATE notes SET notebookId=?,title=?,content=?,contentText=?,contentFormat=?,colorMark=?, isPinned=?,isFavorite=?,isLocked=?,isArchived=?,isTrashed=?,trashedAt=?,version=?,sortOrder=?,updatedAt=? WHERE scopeKey=? AND id=?${encrypted ? " AND version=?" : ""}`, [ @@ -443,11 +443,12 @@ export class NativeLocalRepository implements LocalRepository { scope.scopeKey, id, ...(encrypted ? [current.version] : []), ]); if (encrypted && changed.changes !== 1) throw Object.assign(new Error("笔记版本已改变"), { status: 409, code: "VERSION_CONFLICT" }); - await this.enqueue(tx, "note", id, "upsert", next as unknown as Record, current.version); + await this.enqueue(tx, "note", id, "upsert", next as unknown as Record, current.version, scope.scopeKey); + return { id, savedAt: next.updatedAt }; }); // Refresh per-note receipt from committed Outbox, even while a push is in flight. notifyMobileSyncStatusChanged(); - return { id, savedAt: next.updatedAt }; + return result; } private async removeNote(id: string): Promise { @@ -458,7 +459,7 @@ export class NativeLocalRepository implements LocalRepository { await this.db.transaction(async (tx) => { await tx.run("DELETE FROM notes WHERE scopeKey=? AND id=?", [scope.scopeKey, id]); await tx.run("DELETE FROM native_runtime_meta WHERE key=?", [unsentLocalNoteKey(scope.scopeKey, id)]); - await this.enqueue(tx, "note", id, "delete", undefined, current?.version ?? null); + await this.enqueue(tx, "note", id, "delete", undefined, current?.version ?? null, scope.scopeKey); }); notifyMobileSyncStatusChanged(); } diff --git a/frontend/tests/android-sync/index.html b/frontend/tests/android-sync/index.html new file mode 100644 index 000000000..ae407b276 --- /dev/null +++ b/frontend/tests/android-sync/index.html @@ -0,0 +1,2 @@ +Isolated native sync acceptance +

Initializing isolated native sync acceptance…

diff --git a/frontend/tests/android-sync/main.ts b/frontend/tests/android-sync/main.ts new file mode 100644 index 000000000..540b447c7 --- /dev/null +++ b/frontend/tests/android-sync/main.ts @@ -0,0 +1,59 @@ +// Test-only entry. Included only by vite.android-sync.config.ts, never normal builds. +import { openNativeDatabase } from "../../src/lib/nativeDatabase"; +import { createNativeAttachmentStore } from "../../src/lib/nativeAttachmentStore"; +import { NativeLocalRepository } from "../../src/lib/nativeLocalRepository"; +import { MobileSyncEngine } from "../../src/lib/mobileSyncEngine"; +import { installMobileSyncRealtime } from "../../src/lib/mobileSyncRealtime"; +import { setMobileSyncEnabled } from "../../src/lib/mobileSyncStatus"; +import { readNativeNoteSyncReceipt } from "../../src/lib/mobileNoteSyncReceipt"; +import { realtime } from "../../src/lib/realtime"; + +type Config = { serverUrl: string; token: string; userId: string }; +async function initialize(config: Config) { + localStorage.setItem("sync-acceptance-config", JSON.stringify(config)); + localStorage.setItem("nowen-token", config.token); + localStorage.setItem("nowen-server-url", config.serverUrl); + const profileId = "isolated-acceptance", deviceId = "android-acceptance"; + const accountId = `sync-acceptance:${config.userId}`; + const db = await openNativeDatabase(accountId); + const attachments = await createNativeAttachmentStore(accountId); + await db.run(`INSERT OR IGNORE INTO sync_profiles (id,name,serverUrl,remoteUserId,enabled,createdAt,updatedAt) + VALUES (?,'Acceptance',?,?,1,'now','now')`, [profileId, config.serverUrl, config.userId]); + await db.run("INSERT OR IGNORE INTO sync_devices (profileId,deviceId,platform,createdAt) VALUES (?,?,'android','now')", [profileId, deviceId]); + const engine = new MobileSyncEngine({ db, attachments, profileId, deviceId, ...config }); + const repository = new NativeLocalRepository({ db, attachments, accountId, userId: config.userId, + getScopeKey: () => "personal", requestSync: () => engine.requestSync(150) }); + let dispose: (() => void) | undefined; + const ids = async () => { + const row = (await db.query<{value:string}>("SELECT value FROM native_runtime_meta WHERE key='acceptance-note-ids'"))[0]; + return JSON.parse(row?.value || "[]") as string[]; + }; + const acceptance = { + async createOffline() { + const book = crypto.randomUUID(), notes = [crypto.randomUUID(), crypto.randomUUID()]; + await repository.notebooks.create({ id: book, name: "Acceptance" }); + for (const [index, id] of notes.entries()) await repository.notes.create({ id, notebookId: book, + contentFormat: index ? "tiptap-json" : "markdown", + content: index ? JSON.stringify({ type: "doc", content: [{ type: "paragraph", content: [{ type: "text", text: "offline-rich" }] }] }) : "offline-markdown", + contentText: index ? "offline-rich" : "offline-markdown" }); + await db.run("INSERT INTO native_runtime_meta (key,value,updatedAt) VALUES ('acceptance-note-ids',?,'now') ON CONFLICT(key) DO UPDATE SET value=excluded.value", [JSON.stringify(notes)]); + return notes; + }, + async status() { + return Promise.all((await ids()).map(async (id) => ({ id, note: await repository.notes.get(id), receipt: await readNativeNoteSyncReceipt(db, profileId, id) }))); + }, + async edit(id: string, text: string) { + const note = await repository.notes.get(id); + await repository.notes.update(id, { contentText: text, content: note?.contentFormat === "markdown" ? text : JSON.stringify({ type: "doc", content: [{ type: "paragraph", content: [{ type: "text", text }] }] }) }); + }, + start() { setMobileSyncEnabled(true); dispose = installMobileSyncRealtime(engine); engine.start(); }, + stop() { dispose?.(); dispose = undefined; engine.stop(); setMobileSyncEnabled(false); realtime.disconnect(); }, + async sync() { await engine.syncOnce(); }, + async errors() { return db.query("SELECT lastError FROM sync_state WHERE profileId=?", [profileId]); }, + }; + Object.assign(window, { acceptance, acceptanceRealtime: realtime }); + document.getElementById("status")!.textContent = "Native SQLite ready. Test account only."; +} +Object.assign(window, { initializeAcceptance: initialize }); +const saved = localStorage.getItem("sync-acceptance-config"); +if (saved) void initialize(JSON.parse(saved)).catch(() => { document.getElementById("status")!.textContent = "Native initialization failed"; }); diff --git a/frontend/vite.android-sync.config.ts b/frontend/vite.android-sync.config.ts new file mode 100644 index 000000000..5ff762a28 --- /dev/null +++ b/frontend/vite.android-sync.config.ts @@ -0,0 +1,7 @@ +import path from "node:path"; +import { defineConfig, mergeConfig } from "vite"; +import appConfig from "./vite.config"; +export default mergeConfig(appConfig, defineConfig({ build: { rollupOptions: { input: { + app: path.resolve(__dirname, "index.html"), + acceptance: path.resolve(__dirname, "tests/android-sync/index.html"), +} } } })); diff --git a/scripts/android-sync-acceptance.init.gradle b/scripts/android-sync-acceptance.init.gradle new file mode 100644 index 000000000..5760b8cbb --- /dev/null +++ b/scripts/android-sync-acceptance.init.gradle @@ -0,0 +1,8 @@ +// Explicit debug acceptance build only. Never modifies the production package. +allprojects { + afterEvaluate { p -> + if (p.name == 'app') { + p.android.defaultConfig.applicationId = 'com.nowen.note.syncacceptance' + } + } +} diff --git a/scripts/android-sync-acceptance.mjs b/scripts/android-sync-acceptance.mjs new file mode 100644 index 000000000..59c2e37c0 --- /dev/null +++ b/scripts/android-sync-acceptance.mjs @@ -0,0 +1,105 @@ +// Runs only against the separate com.nowen.note.syncacceptance debug APK. +// Requires backend/tests/sync-v2-android-server.ts on loopback port 47831. +import assert from "node:assert/strict"; +import { execFileSync } from "node:child_process"; +import { createRequire } from "node:module"; +import fs from "node:fs"; +import path from "node:path"; +const require = createRequire(new URL("../frontend/package.json", import.meta.url)); +const { _android } = require("@playwright/test"); +const serial = process.argv[2]; +if (!serial) throw new Error("Pass the isolated test device serial"); +const adbPath = process.env.ADB || path.join(process.env.LOCALAPPDATA, "Android/Sdk/platform-tools/adb.exe"); +const adb = (...args) => execFileSync(adbPath, ["-s", serial, ...args], { encoding: "utf8", timeout: 30_000 }); +const base = "http://127.0.0.1:47831"; +const config = await (await fetch(`${base}/acceptance/config`)).json(); +let browser, page; +async function attach() { + const devices = await _android.devices(); + browser = devices.find((device) => device.serial() === serial); + assert.ok(browser, "isolated acceptance device is available"); + page = await (await browser.webView({ pkg: "com.nowen.note.syncacceptance" })).page(); + await page.goto("https://localhost/tests/android-sync/index.html"); + await page.waitForFunction(() => typeof window.initializeAcceptance === "function"); + await page.evaluate(async (config) => { + await window.initializeAcceptance(config); + }, config); + await page.waitForFunction(() => !!window.acceptance, { timeout: 20_000 }); +} +async function wait(work, message, timeout = 12_000) { + const end = Date.now() + timeout; + while (Date.now() < end) { + if (await work()) return; + await new Promise((resolve) => setTimeout(resolve, 25)); + } + throw new Error(message); +} +const native = (method, ...args) => page.evaluate(({ method, args }) => window.acceptance[method](...args), { method, args }); +const serverNote = async (id) => (await fetch(`${base}/acceptance/note/${id}`)).json(); +const remoteEdit = async (id, text) => { + const response = await fetch(`${base}/acceptance/edit/${id}`, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify({ text }) }); + assert.ok(response.ok); + const body = await response.json(); + assert.equal(body.results[0].status, "applied"); +}; +const percentile = (values, p) => [...values].sort((a, b) => a - b)[Math.ceil(p * values.length) - 1]; +const stats = (values) => ({ samples: values.length, p50Ms: percentile(values, .5), p95Ms: percentile(values, .95), p99Ms: percentile(values, .99) }); +try { + await attach(); + adb("reverse", "--remove", "tcp:47831"); + const unreachable = await page.evaluate(async () => { + try { await fetch("http://127.0.0.1:47831/acceptance/config"); return false; } + catch { return true; } + }); + assert.ok(unreachable, "Android cannot reach the test server during offline edits"); + await native("start"); + const ids = await native("createOffline"); + let rows = await native("status"); + assert.deepEqual(rows.map((r) => r.note.contentText), ["offline-markdown", "offline-rich"]); + assert.ok(rows.every((r) => r.receipt.phase === "pending")); + await browser.close(); + adb("shell", "am", "force-stop", "com.nowen.note.syncacceptance"); + adb("shell", "am", "start", "-n", "com.nowen.note.syncacceptance/com.nowen.note.MainActivity"); + await attach(); + rows = await native("status"); + assert.deepEqual(rows.map((r) => r.note.id), ids); + assert.deepEqual(rows.map((r) => r.note.contentText), ["offline-markdown", "offline-rich"]); + assert.ok(rows.every((r) => r.receipt.phase === "pending")); + adb("reverse", "tcp:47831", "tcp:47831"); + await native("start"); + await wait(async () => (await native("status")).every((r) => r.receipt.phase === "confirmed"), "offline edits not acknowledged"); + const download = [], upload = []; + for (let i = 0; i < 30; i++) { + const id = ids[i % 2], remoteText = `remote-${i}`, localText = `native-${i}`; + let started = performance.now(); + await remoteEdit(id, remoteText); + await wait(async () => (await native("status")).find((r) => r.id === id).note.contentText === remoteText, "remote update not pulled"); + download.push(Math.round(performance.now() - started)); + started = performance.now(); + await native("edit", id, localText); + await wait(async () => (await serverNote(id))?.contentText === localText, "native update not pushed"); + await wait(async () => (await native("status")).find((r) => r.id === id).receipt.phase === "confirmed", "edit ACK not persisted"); + upload.push(Math.round(performance.now() - started)); + } + // Remove the notification transport while leaving the same engine's poll active. + await page.evaluate(async () => { const realtime = window.acceptanceRealtime; realtime.disconnect(); }); + await remoteEdit(ids[0], "lost-notice"); + await wait(async () => (await native("status"))[0].note.contentText === "lost-notice", "poll did not compensate for missing notification", 40_000); + await native("stop"); + await remoteEdit(ids[0], "reconnected"); + await native("start"); + await wait(async () => (await native("status"))[0].note.contentText === "reconnected", "reconnect did not catch up"); + const report = { device: serial, package: "com.nowen.note.syncacceptance", offlineImmediateRead: true, forceStopRecovery: true, + offlineNetworkIsolation: "adb reverse removed; HTTP server unreachable from Android", + markdownAndRichText: true, missingNoticePoll: true, reconnect: true, + httpPushToNativePull: stats(download), nativeCommitToServerAndAck: stats(upload), errors: await native("errors"), + boundary: "Loopback adb reverse; real Android native SQLite and MobileSyncEngine. Does not test full EditorPane, knowledge tree creation, attachment bytes, or production JWT middleware." }; + fs.mkdirSync(new URL("../docs/test-results/", import.meta.url), { recursive: true }); + fs.writeFileSync(new URL("../docs/test-results/sync-v2-android-2026-10-11.json", import.meta.url), JSON.stringify(report, null, 2)); + console.log(JSON.stringify(report, null, 2)); +} finally { + if (page) await native("stop").catch(() => undefined); + await browser?.close().catch(() => undefined); + adb("reverse", "--remove", "tcp:47831"); +} +process.exit(0); diff --git a/scripts/data-consistency-contract.mjs b/scripts/data-consistency-contract.mjs index 4f02d05ff..5dda41548 100644 --- a/scripts/data-consistency-contract.mjs +++ b/scripts/data-consistency-contract.mjs @@ -46,7 +46,7 @@ if ([...argumentsSet].some((arg) => !supported.includes(arg)) || argumentsSet.si } for (const check of policy.checks) { if (typeof check.file !== "string" || typeof check.case !== "string" || - check.case.length < 12 || !/^(frontend\/src\/lib\/__tests__|backend\/tests)\/[\w./-]+\.test\.tsx?$/.test(check.file)) { + check.case.length < 12 || !/^(frontend\/src\/lib\/(?:__tests__|encryptedNotes\/__tests__)|backend\/tests)\/[\w./-]+\.test\.tsx?$/.test(check.file)) { failures.push(policy.id + " uses an invalid test path/case: " + JSON.stringify(check)); continue; }