diff --git a/docs/design.md b/docs/design.md index 5b31767..8b292ff 100644 --- a/docs/design.md +++ b/docs/design.md @@ -465,7 +465,7 @@ crsl-libはMonasのために設計されたCIDネイティブなDAG CRDTライ | フィールド | 規則 | 理由 | |---|---|---| -| コンテンツ本体(ciphertext) | LWW — ただし**親から本体を変えたheadの中で**timestamp最大 | 最後の書き込みが正。policyだけを進めたhead(revoke等)は本体を親からコピーしているだけで「書いて」いないので、並行する本物の書き込みに勝ってはならない | +| コンテンツ本体(ciphertext) | `(body_updated_at, data)` の辞書順 max | 明示的な本文更新でのみ順序を進め、policy-only 更新と Merge は本文と順序をそのまま引き継ぐ。head 自体の timestamp や直近の親との差分では選ばない | | `access_policy.min_valid_issued_at` | 全headの**max** | 失効境界は単調にしか進まない。timestampで選ぶと、境界を知らないノードが受理した並行writeが境界を巻き戻す | | `access_policy.owner` / `content_id` | 不変(genesisで確定) | — | @@ -473,6 +473,12 @@ payload全体をtimestampで丸ごと選ぶ(純粋なLWW)と、revokeと並 このマージが決めるのは「Mergeノードに何を入れるか」であり、「そのheadを受理してよかったか」ではない。失効境界を知らないメンバーが旧Tokenで受理したwriteは、最新のwriteであれば本体として残る(境界は残るので以後は書けない)。それを弾くにはwriteが自分のTokenを持ち歩き、マージ時に畳んだ境界に対して検証する必要がある — ワイヤ形式の変更を伴うため別issueで追跡する。 +`body_updated_at` は本文と同じ payload に保存する論理的な更新順序であり、別 DAG ではない。本文更新ではローカルの単調 timestamp と観測済みの順序 + 1 の大きい方を採る。policy-only 更新、再マージ、再起動で順序を失わず、同値時は ciphertext の辞書順で決定する。観測・マージ・payload の生成・commit は同じ repository lock 内で行う。 + +同期 export は operation と DAG ノードを payload・parents・genesis・metadata で対応付け、実ノードの timestamp を送る。履歴の位置対応は使用しない。`since_version` はそのノードと祖先を既知とみなし、兄弟枝を省かず親から順に送る。曖昧な対応は推測せずエラーにする。 + +保存・wire 形式の変更: `body_updated_at` は必須で、旧形式を 0 等へ暗黙補完しない。現行デモは顧客利用前のため、全 state-node を同時更新し、新しいストアから開始してコンテンツを再作成する必要がある。既存ストアを維持する場合の移行は未実装。データ削除やデプロイは本変更では行わない。 + コンテンツ本体の意味的なマージ(同じフィールド内での両立)は現時点で未実装であり、研究課題として位置づけられている。 ### 将来のCRDT拡張 diff --git a/example-ui/README.md b/example-ui/README.md index fd642b8..98dfdbb 100644 --- a/example-ui/README.md +++ b/example-ui/README.md @@ -51,6 +51,32 @@ send permissive CORS headers. ## Tests +### Isolated UI regressions (no backend required) + +```bash +npm ci --ignore-scripts +npx playwright install chromium # once, if not already installed +npm run test:regression +npm run build +``` + +This suite starts its own Vite on `127.0.0.1:5198` (fails if occupied). It runs +actual App/store/flow/API-adapter paths with HTTP fixtures, intercepts all +backend requests, rejects unknown endpoints and blocks external origins. No +hosted nodes or account keys are needed or modified. Results/traces go to +`/tmp/monas-ui-regression-results`. It covers recipient import → edit → reopen, +owner preview/head checks, revocation reach reporting, and legacy identities. +These are UI regressions, not cryptographic or distributed-protocol tests. + +Legacy identity migration keeps the **last-created signing account**, matching +`POST /accounts` replacing monas-account's single key. Earlier signing entries +remain available as envelope-decryption keypairs; removing the current account +does not promote them. `activeLabel` from old UI switching cannot change the +backend's key. The account API has no read-current-key endpoint, so a reset or +externally replaced backend key still requires explicit user recovery. + +### Real-stack suites + ```bash npm test # UI suite (tests/) — ~22s npm run test:ui # same, in the Playwright UI runner diff --git a/example-ui/package.json b/example-ui/package.json index 0ea3fee..f3cd68c 100644 --- a/example-ui/package.json +++ b/example-ui/package.json @@ -10,6 +10,7 @@ "preview": "vite preview", "test": "playwright test", "test:ui": "playwright test --ui", + "test:regression": "playwright test -c playwright.regression.config.ts", "test:e2e": "playwright test -c playwright.e2e.config.ts" }, "dependencies": { diff --git a/example-ui/playwright.regression.config.ts b/example-ui/playwright.regression.config.ts new file mode 100644 index 0000000..5911cb1 --- /dev/null +++ b/example-ui/playwright.regression.config.ts @@ -0,0 +1,18 @@ +import { defineConfig, devices } from "@playwright/test"; + +// No running gateway/account/state nodes required. Never reuse a live UI server. +export default defineConfig({ + testDir: "./tests-regression", + timeout: 30_000, + expect: { timeout: 4_000 }, + workers: 1, + reporter: "list", + outputDir: "/tmp/monas-ui-regression-results", + use: { baseURL: "http://127.0.0.1:5198", serviceWorkers: "block", trace: "retain-on-failure" }, + webServer: { + command: "npx vite --host 127.0.0.1 --port 5198 --strictPort", + url: "http://127.0.0.1:5198", + reuseExistingServer: false, + }, + projects: [{ name: "chromium", use: { ...devices["Desktop Chrome"] } }], +}); diff --git a/example-ui/src/App.tsx b/example-ui/src/App.tsx index 2b4d892..438aefd 100644 --- a/example-ui/src/App.tsx +++ b/example-ui/src/App.tsx @@ -579,7 +579,12 @@ export default function App() { // voided tokens until its next sync. Say so — "revoked" alone would // overstate what just happened. const reach = r?.token_invalidation_reach; - if (reach?.relayed) { + if (!reach && (r?.token_invalidated_at != null || entry.syncedToStateNode || entry.remoteContentId)) { + pushToast( + "Token invalidation propagation is unknown — the node returned no reach report. Members that have not synced may still accept writes under old tokens.", + "error", + ); + } else if (reach?.relayed) { pushToast( "The revoke was relayed to a member node; which members enforce the cutoff yet is not known from here. Writes under the old token may land on members that have not synced.", "error", @@ -806,6 +811,7 @@ export default function App() { handleEditOpen(e)} onClose={() => setModal({ type: "none" })} diff --git a/example-ui/src/api/share.ts b/example-ui/src/api/share.ts index 39b8388..ba1d58b 100644 --- a/example-ui/src/api/share.ts +++ b/example-ui/src/api/share.ts @@ -88,7 +88,8 @@ export interface RevokeShareOutput { * keeps accepting writes under the voided tokens until its next sync. * The revoke does not wait for that (a writer must not be able to block * it); this is how the UI tells "revoked everywhere" from "revoked, N - * members still to hear". Absent when no state node was involved. */ + * members still to hear". Also absent with legacy nodes (unknown reach); + * absence alone does not mean no state node was involved. */ token_invalidation_reach?: TokenInvalidationReach; } diff --git a/example-ui/src/components/IdentityModal.tsx b/example-ui/src/components/IdentityModal.tsx index 537d9cb..5eeb837 100644 --- a/example-ui/src/components/IdentityModal.tsx +++ b/example-ui/src/components/IdentityModal.tsx @@ -78,7 +78,7 @@ export function IdentityModal({ onClose }: { onClose: () => void }) { signing ) : ( - + keypair only )} @@ -142,7 +142,7 @@ export function IdentityModal({ onClose }: { onClose: () => void }) { Other identities
- Keypair-only identities from an earlier version of this UI. They can + Older identities (including replaced signing accounts). They can still open envelopes addressed to them, but cannot sign or read the state node. Have files shared to your account instead.
diff --git a/example-ui/src/components/PreviewModal.tsx b/example-ui/src/components/PreviewModal.tsx index e344a44..1870eef 100644 --- a/example-ui/src/components/PreviewModal.tsx +++ b/example-ui/src/components/PreviewModal.tsx @@ -45,12 +45,15 @@ function classifyVerify(v: VerifyIntegrityOutput): "valid" | "behind" | "no-loca export function PreviewModal({ entry, contentB64Url, + displayedVersionId, onCheckHead, onEdit, onClose, }: { entry: Entry; contentB64Url: string; + /** Identity of the rendered body, fixed when opened (not the last write). */ + displayedVersionId: string | undefined; /** Verified read of the head that also records it on the entry — the * list's sync badge and the status line here read from that record. */ onCheckHead: (entry: Entry) => Promise; @@ -166,7 +169,7 @@ export function PreviewModal({ // The one-line answer to "am I looking at the newest version?". Derived // from what the entry recorded at the last head check (open, sweep, or the // verified read below), so it stays right even while this dialog is idle. - const sync = syncStatusOf(entry); + const sync = syncStatusOf(entry, displayedVersionId); const syncText = describeSync(sync); const held = heldVersionId(entry); const syncBadgeClass = @@ -227,7 +230,7 @@ export function PreviewModal({ <> The Content Network has a newer version than the text above {received - ? " — the owner has edited since sharing" + ? " — this preview is the original shared version, not the current network head" : " — a recipient with write access has edited since your last save"} . Read from state-node shows it {canWrite ? (received ? "; Edit contents starts from it." : "; Pull & edit adopts it into your copy.") : "."} diff --git a/example-ui/src/pipeline/flows.ts b/example-ui/src/pipeline/flows.ts index 011f11e..61b8cd9 100644 --- a/example-ui/src/pipeline/flows.ts +++ b/example-ui/src/pipeline/flows.ts @@ -485,7 +485,11 @@ export function revokeFlow(input: { exec: async (ctx) => { const r = ctx.revoke as shareApi.RevokeShareOutput; const reach = r.token_invalidation_reach; - if (!reach) return "No state node involved — nothing to propagate"; + if (!reach) { + return r.token_invalidated_at != null || entry.syncedToStateNode || entry.remoteContentId + ? "Token invalidation propagation is unknown — the node returned no reach report. Members that have not synced may still accept writes under old tokens." + : "No state node involved — nothing to propagate"; + } if (reach.relayed) { return ( "The node we contacted relayed the revoke to a member; which members " + diff --git a/example-ui/src/store/identity.ts b/example-ui/src/store/identity.ts index 07543fa..1020fa6 100644 --- a/example-ui/src/store/identity.ts +++ b/example-ui/src/store/identity.ts @@ -18,6 +18,23 @@ const store = createStore("monas.identities.v2", { activeLabel: null, }); +// Older UIs appended a signing account on every POST /accounts, but that +// endpoint replaces the backend's ONE key. Array order records creation order; +// activeLabel only recorded UI switching and cannot change the backend key. +// There is no account read endpoint to reconcile against. Retain old private +// keys for envelope decryption, but persist their demotion so removing the +// current account never resurrects an overwritten signing key. +const signingAccounts = store.get().identities.filter((i) => i.isSigningAccount); +if (signingAccounts.length > 1) { + const current = signingAccounts[signingAccounts.length - 1]; + store.set((prev) => ({ + identities: prev.identities.map((i) => + i.isSigningAccount && i !== current ? { ...i, isSigningAccount: false } : i, + ), + activeLabel: current.label, + })); +} + export const useIdentities = () => store.use(); export function getIdentities(): Identity[] { diff --git a/example-ui/src/store/sync.ts b/example-ui/src/store/sync.ts index d6aa014..0967906 100644 --- a/example-ui/src/store/sync.ts +++ b/example-ui/src/store/sync.ts @@ -31,7 +31,7 @@ export function heldVersionId(entry: Entry): string | undefined { const checking = new Set(); -export function syncStatusOf(entry: Entry): SyncStatus { +export function syncStatusOf(entry: Entry, comparedVersionId = heldVersionId(entry)): SyncStatus { if (entry.kind !== "file" || !entry.syncedToStateNode || !entry.remoteContentId) { return { kind: "local" }; } @@ -40,7 +40,7 @@ export function syncStatusOf(entry: Entry): SyncStatus { return { kind: "unreachable", error: entry.networkCheckError, head: entry.networkHead }; } if (!entry.networkHead) return { kind: "unchecked" }; - return entry.networkHead.localId === heldVersionId(entry) + return entry.networkHead.localId === comparedVersionId ? { kind: "current", head: entry.networkHead } : { kind: "behind", head: entry.networkHead }; } diff --git a/example-ui/tests-regression/review.spec.ts b/example-ui/tests-regression/review.spec.ts new file mode 100644 index 0000000..aee381a --- /dev/null +++ b/example-ui/tests-regression/review.spec.ts @@ -0,0 +1,180 @@ +import { test, expect, type Page } from "@playwright/test"; +import type { Entry, Identity, SharePackage } from "../src/types"; + +const identity = (label: string, signing = true): Identity => ({ label, keyType: "secp256r1", publicKeyB64Url: `public-${label}`, privateKeyB64Url: `fake-private-${label}`, isSigningAccount: signing }); +const envelope = { enc: "fixture", wrapped_cek: "fixture", ciphertext: "fixture-v1", key_epoch: 1 }; +const pkg: SharePackage = { + kind: "monas-share", v: 1, name: "shared.txt", mimeType: "text/plain", sizeBytes: 2, + content_id: "plain-v1", remote_content_id: "network", sender_public_key: "public-owner", + recipient_public_key: "public-B", recipient_key_id: "recipient-B", permissions: ["read", "write"], + key_envelope: envelope, shared_at: 1, + delegated_access: { delegated_token: "fixture-token", issued_at: 1, expires_at: 4102444800, jti: "fixture" }, +}; +const ownerEntry = (): Entry => ({ id: "owner-file", kind: "file", name: "owner.txt", parentPath: "/", sizeBytes: 2, mimeType: "text/plain", createdAt: 1, updatedAt: 1, localContentId: "plain-v1", remoteContentId: "network", syncedToStateNode: true, versionCount: 1, shares: [] }); +const b64 = (text: string) => Buffer.from(text).toString("base64url"); + +// Exercise App, stores, flow runner and API adapters unchanged. Only HTTP is +// simulated: every gateway call is intercepted, unknown calls fail closed. +async function boot(page: Page, entries: Entry[] = [], identities = [identity("B")], activeLabel = "B") { + let head = "v1"; + const requests: { path: string; body: any; method: string }[] = []; + const unexpected: string[] = []; + await page.addInitScript(({ entries, identities, activeLabel }) => { + if (sessionStorage.getItem("regression-seeded")) return; + sessionStorage.setItem("regression-seeded", "yes"); + localStorage.setItem("monas.identities.v2", JSON.stringify({ identities, activeLabel })); + localStorage.setItem("monas.registry.v3", JSON.stringify(entries)); + localStorage.setItem("monas.endpoints.v2", JSON.stringify({ gateway: "/api", accountService: "/account-api" })); + }, { entries, identities, activeLabel }); + await page.route("**/*", async route => { + const req = route.request(); + const url = new URL(req.url()); + const path = url.pathname; + if (url.origin !== "http://127.0.0.1:5198") { + unexpected.push(req.url()); + return route.abort(); + } + if (!path.startsWith("/api") && !path.startsWith("/account-api")) return route.continue(); + const body = req.postDataJSON(); + requests.push({ path, body, method: req.method() }); + let data: unknown; + if (path === "/api/health") return route.fulfill({ status: 200, body: "" }); + if (path === "/api/share/decrypt") data = { content_id: "plain-v1", content: b64("V1"), version: "v1" }; + else if (path === "/api/state/read") data = { content_id: "network", local_content_id: `plain-${head}`, version: head, content: b64(head.toUpperCase()) }; + else if (path === "/api/state/latest-version") data = { content_id: "network", latest_version: head }; + else if (path === "/api/state/history") data = { content_id: "network", versions: [head] }; + else if (path === "/api/share/content/network" && req.method() === "PUT") { + expect(body.content).toBe(b64("V2")); + expect(req.headers().authorization).toBe("Bearer fixture-token"); + head = "v2"; + data = { remote_content_id: "network", version_id: "plain-v2" }; + } else if (path === "/api/content/plain-v1") data = { content_id: "plain-v1", content: b64("V1") }; + else { + unexpected.push(`${req.method()} ${path}`); + return route.fulfill({ status: 500, body: "Unexpected test request" }); + } + return route.fulfill({ json: { success: true, data, trace_id: "test" } }); + }); + await page.goto("/"); + return { requests, unexpected, setHead: (value: string) => { head = value; } }; +} +async function action(page: Page, name: string) { + await page.locator(".row-menu-wrap > button").click(); + await page.locator(".menu").getByRole("button", { name, exact: true }).click(); +} + +test("recipient import → edit → reopen compares the displayed envelope, not the recorded write", async ({ page }) => { + const gateway = await boot(page); + await page.getByRole("button", { name: "Import shared", exact: true }).click(); + await page.locator(".modal textarea").fill(JSON.stringify(pkg)); + await page.getByRole("button", { name: "Unwrap & add to my Drive" }).click(); + await expect(page.locator(".preview-box").first()).toHaveText("V1"); + await expect(page.locator(".sync-status")).toHaveAttribute("data-sync", "current"); + await page.keyboard.press("Escape"); + await action(page, "Edit contents"); + await expect(page.locator(".modal textarea")).toHaveValue("V1"); + await page.locator(".modal textarea").fill("V2"); + await page.getByRole("button", { name: "Re-encrypt & save" }).click(); + await expect.poll(() => page.evaluate(() => JSON.parse(localStorage.getItem("monas.registry.v3")!)[0].receivedShare.writtenVersionId)).toBe("plain-v2"); + await expect(page.locator(".row .sync")).toHaveAttribute("data-sync", "current"); + await action(page, "Open / preview"); + await expect(page.locator(".preview-box").first()).toHaveText("V1"); + await expect(page.locator(".sync-status")).toHaveAttribute("data-sync", "behind"); + await expect(page.locator(".sync-status")).not.toContainText("the owner has edited"); + await page.getByRole("button", { name: "Read from state-node", exact: true }).click(); + await expect(page.locator(".preview-box").last()).toHaveText("V2"); + await expect(page.locator(".modal")).toContainText("your edit"); + await expect(page.locator(".sync-status")).toHaveAttribute("data-sync", "behind"); + await expect(page.locator(".row .sync")).toHaveAttribute("data-sync", "current"); + expect(gateway.requests.filter(r => r.path === "/api/share/decrypt")).toHaveLength(2); + expect(gateway.unexpected).toEqual([]); +}); + +test("owner preview compares its local body with the network head", async ({ page }) => { + const gateway = await boot(page, [ownerEntry()]); + await action(page, "Open / preview"); + await expect(page.locator(".preview-box").first()).toHaveText("V1"); + await expect(page.locator(".sync-status")).toHaveAttribute("data-sync", "current"); + gateway.setHead("v2"); + await page.getByRole("button", { name: "Check now" }).click(); + await expect(page.locator(".sync-status")).toHaveAttribute("data-sync", "behind"); + await expect(page.locator(".row .sync")).toHaveAttribute("data-sync", "behind"); + await expect(page.locator(".preview-box").first()).toHaveText("V1"); + await expect(page.locator(".preview-box").last()).toHaveText("V2"); + expect(gateway.unexpected).toEqual([]); +}); + +const reachCases = [ + { name: "legacy", output: { token_invalidated_at: 42 }, detail: /propagation.*unknown/i, toast: /propagation.*unknown/i }, + { name: "local", output: {}, detail: /No state node involved/, toast: null }, + { name: "relayed", output: { token_invalidated_at: 42, token_invalidation_reach: { notified_members: [], unreached_members: [], relayed: true } }, detail: /relayed.*not known/, toast: /relayed to a member/ }, + { name: "all", output: { token_invalidated_at: 42, token_invalidation_reach: { notified_members: ["node-2"], unreached_members: [], relayed: false } }, detail: /All 1 other member/, toast: null }, + { name: "unreached", output: { token_invalidated_at: 42, token_invalidation_reach: { notified_members: [], unreached_members: [{ node_id: "node-2", error: "offline" }], relayed: false } }, detail: /1 member\(s\) did not get/, toast: /did not reach 1 member/ }, +]; +for (const scenario of reachCases) { + test(`revoke ${scenario.name}: flow and toast report propagation honestly`, async ({ page }) => { + const entry = ownerEntry(); + if (scenario.name === "local") { entry.syncedToStateNode = false; delete entry.remoteContentId; } + entry.shares = [{ recipientPublicKeyB64Url: "public-recipient", recipientLabel: "recipient", permissions: ["read", "write"], senderKeyId: "B", recipientKeyId: "recipient", senderPublicKeyB64Url: "public-B", envelope, grantedAt: 1 }]; + const gateway = await boot(page, [entry]); + await page.route("**/api/share/revoke", async route => { + expect(route.request().postDataJSON().sender_public_key).toBe("public-B"); + await route.fulfill({ json: { success: true, data: { content_id: "plain-v1", recipient_public_key: "public-recipient", revoked: true, reissued_envelopes: [], ...scenario.output } } }); + }); + await action(page, "Share"); + await page.getByRole("button", { name: "Revoke", exact: true }).click(); + await expect(page.locator(".run-status")).toHaveText("complete"); + // Soft assertions catch flow and App toast independently in a RED run. + await expect.soft(page.locator(".step").filter({ hasText: "Cutoff propagation" })).toContainText(scenario.detail); + if (scenario.toast) await expect.soft(page.locator(".toast.error")).toContainText(scenario.toast); + else await expect(page.locator(".toast.error")).toHaveCount(0); + await expect(page.locator(".modal").getByRole("button", { name: "Revoke", exact: true })).toHaveCount(0); + expect(gateway.unexpected).toEqual([]); + }); +} + +for (const activeLabel of ["B", "A", "legacy"]) { + test(`identity migration retains latest signing key B when old activeLabel is ${activeLabel}`, async ({ page }) => { + const ids = [identity("A"), identity("B"), identity("legacy", false)]; + const gateway = await boot(page, [], ids, activeLabel); + await expect.soft(page.locator(".account-chip .meta b")).toHaveText("B"); + await page.locator(".account-chip").click(); + const signing = page.locator(".recipient-row").filter({ has: page.locator(".badge.enc") }); + await expect.soft(signing).toContainText("B"); + await expect.soft(page.locator(".recipient-row")).toHaveCount(3); + await expect.soft(page.locator(".recipient-row").filter({ hasText: "public-A" })).toContainText("keypair only"); + await expect.soft(page.locator(".recipient-row").filter({ hasText: "public-legacy" })).toContainText("keypair only"); + const migrated = await page.evaluate(() => JSON.parse(localStorage.getItem("monas.identities.v2")!).identities); + expect(migrated).toEqual([{ ...ids[0], isSigningAccount: false }, ids[1], ids[2]]); + await page.reload(); + await expect(page.locator(".account-chip .meta b")).toHaveText("B"); + await page.locator(".account-chip").click(); + // Migration preserves all material, but never revives A as signing when B is removed. + await signing.getByTitle("Remove from this browser").click(); + await expect(page.getByRole("button", { name: "Create account", exact: true })).toBeVisible(); + await expect(page.locator(".recipient-row .badge.enc")).toHaveCount(0); + const stored = await page.evaluate(() => JSON.parse(localStorage.getItem("monas.identities.v2")!)); + expect(stored.identities.every((i: Identity) => !i.isSigningAccount)).toBe(true); + expect(stored.identities.find((i: Identity) => i.label === "legacy")).toEqual(ids[2]); + expect(gateway.unexpected).toEqual([]); + }); +} + +test("identity with one signing key preserves it over an active legacy keypair", async ({ page }) => { + await boot(page, [], [identity("old", false), identity("B"), identity("legacy", false)], "legacy"); + await expect(page.locator(".account-chip .meta b")).toHaveText("B"); + await page.locator(".account-chip").click(); + await expect(page.locator(".recipient-row")).toHaveCount(3); + await expect(page.locator(".recipient-row").filter({ has: page.locator(".badge.enc") })).toContainText("B"); +}); + +test("identity keypair-only fallback preserves active legacy identity and all envelope keys", async ({ page }) => { + const ids = [identity("A", false), identity("B", false)]; + await boot(page, [], ids, "B"); + await expect(page.locator(".account-chip .meta b")).toHaveText("B"); + await page.locator(".account-chip").click(); + await expect(page.locator(".recipient-row")).toHaveCount(2); + await expect(page.locator(".recipient-row .badge.enc")).toHaveCount(0); + await expect(page.getByRole("button", { name: "Create account", exact: true })).toBeVisible(); + expect(await page.evaluate(() => JSON.parse(localStorage.getItem("monas.identities.v2")!).identities)).toEqual(ids); +}); diff --git a/monas-content/src/infrastructure/node_verification.rs b/monas-content/src/infrastructure/node_verification.rs index 0b1dc7a..196dfc5 100644 --- a/monas-content/src/infrastructure/node_verification.rs +++ b/monas-content/src/infrastructure/node_verification.rs @@ -95,7 +95,8 @@ pub fn verify_and_extract( } /// Pull `data: Vec` out of the decoded payload value. The payload is -/// `ContentPayload { data, access_policy }`, CBOR-encoded as a map. +/// `ContentPayload { data, body_updated_at, access_policy }`, CBOR-encoded as +/// a map. Additional payload fields remain covered by the raw-byte CID check. fn extract_payload_data(payload: &serde_cbor::Value) -> Result, NodeVerificationError> { use serde_cbor::Value; match payload { @@ -177,6 +178,37 @@ mod tests { assert!(verified.parents.is_empty()); } + #[test] + fn body_write_order_is_verified_without_affecting_ciphertext_extraction() { + use crsl_lib::convergence::metadata::ContentMetadata; + use crsl_lib::dasl::node::Node; + #[derive(Clone, serde::Serialize, serde::Deserialize)] + struct Payload { + data: Vec, + body_updated_at: u64, + access_policy: Option<()>, + } + let mut node = Node::new_genesis( + Payload { + data: b"ciphertext".to_vec(), + body_updated_at: 10, + access_policy: None, + }, + 20, + ContentMetadata::default(), + ); + let version = node.content_id().unwrap().to_string(); + let verified = verify_and_extract(&node.to_bytes().unwrap(), &version).unwrap(); + assert_eq!(verified.ciphertext, b"ciphertext"); + // Same ciphertext, but changing its ordering metadata must invalidate + // the version CID too. The new field is not outside the trust check. + node.payload.body_updated_at = 11; + assert!(matches!( + verify_and_extract(&node.to_bytes().unwrap(), &version), + Err(NodeVerificationError::CidMismatch { .. }) + )); + } + #[test] fn verify_rejects_tampered_payload() { let bytes = make_node(b"original".to_vec(), vec![]); diff --git a/monas-state-node/src/infrastructure/crdt_repository.rs b/monas-state-node/src/infrastructure/crdt_repository.rs index 8b061d0..bf44046 100644 --- a/monas-state-node/src/infrastructure/crdt_repository.rs +++ b/monas-state-node/src/infrastructure/crdt_repository.rs @@ -16,6 +16,7 @@ use crsl_lib::convergence::metadata::ContentMetadata; use crsl_lib::crdt::crdt_state::CrdtState; use crsl_lib::crdt::operation::{Operation, OperationType}; use crsl_lib::crdt::storage::LeveldbStorage; +use crsl_lib::crdt::timestamp::next_monotonic_timestamp; use crsl_lib::graph::dag::DagGraph; use crsl_lib::graph::storage::{LeveldbNodeStorage, NodeStorage}; use crsl_lib::repo::Repo; @@ -30,6 +31,11 @@ use std::path::Path; #[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] pub struct ContentPayload { pub data: Vec, + /// Ordering of the last explicit body write, not of the containing node. + /// Policy-only updates and merges MUST copy this together with `data`. + /// Required on the wire: silently defaulting old data would invent a write + /// order. Deploy this format to all members together with fresh stores. + pub body_updated_at: u64, pub access_policy: Option, } @@ -171,6 +177,7 @@ impl ContentRepository for CrslCrdtRepository { let placeholder = Self::generate_placeholder_cid(data); let payload = ContentPayload { data: data.to_vec(), + body_updated_at: next_monotonic_timestamp(), access_policy, }; @@ -202,33 +209,36 @@ impl ContentRepository for CrslCrdtRepository { ) -> Result { let genesis = Self::parse_cid(genesis_cid)?; - // If no access_policy provided, preserve the existing one from the latest version - let policy = if access_policy.is_some() { - access_policy - } else { - let mut repo = self.repo.lock(); - Self::converged_head(&mut repo, &genesis)?.and_then(|latest_cid| { - repo.dag - .get_node(&latest_cid) - .ok() - .flatten() - .and_then(|node| node.payload().access_policy.clone()) - }) - }; - + // Read/derive/commit under one lock: a policy arriving between the + // read and commit must not be overwritten by the policy we copied. + let mut repo = self.repo.lock(); + let head = Self::converged_head(&mut repo, &genesis)? + .ok_or_else(|| anyhow::anyhow!("Content not found: {genesis}"))?; + let current = repo + .dag + .get_node(&head)? + .ok_or_else(|| anyhow::anyhow!("Head node not found: {head}"))?; + let current = current.payload(); + let mut policy = access_policy.or_else(|| current.access_policy.clone()); + if let (Some(policy), Some(previous)) = (&mut policy, ¤t.access_policy) { + policy.raise_min_valid_issued_at(previous.min_valid_issued_at()); + } + // Advance past an observed write even if this replica's wall clock + // lags. Equal timestamps on independent replicas are broken by bytes + // in MonasMergePolicy, never by storage/arrival order. + let after_current = current + .body_updated_at + .checked_add(1) + .ok_or_else(|| anyhow::anyhow!("Body write timestamp exhausted"))?; let payload = ContentPayload { data: data.to_vec(), + body_updated_at: next_monotonic_timestamp().max(after_current), access_policy: policy, }; - - // Create update operation - parents will be auto-filled by crsl-lib let op = Operation::new(genesis, OperationType::Update(payload), author.to_string()); - - let version_cid = { - let mut repo = self.repo.lock(); - repo.commit_operation(op) - .map_err(|e| anyhow::anyhow!("Failed to commit update operation: {}", e))? - }; + let version_cid = repo + .commit_operation(op) + .map_err(|e| anyhow::anyhow!("Failed to commit update operation: {}", e))?; Ok(CommitResult { genesis_cid: genesis_cid.to_string(), @@ -376,32 +386,24 @@ impl ContentRepository for CrslCrdtRepository { ) -> Result { let genesis = Self::parse_cid(genesis_cid)?; - // Get current data from latest version - let current_data = { - let mut repo = self.repo.lock(); - Self::converged_head(&mut repo, &genesis)? - .and_then(|latest_cid| { - repo.dag - .get_node(&latest_cid) - .ok() - .flatten() - .map(|node| node.payload().data.clone()) - }) - .unwrap_or_default() - }; - - let payload = ContentPayload { - data: current_data, - access_policy: Some(policy), - }; - + let mut repo = self.repo.lock(); + let head = Self::converged_head(&mut repo, &genesis)? + .ok_or_else(|| anyhow::anyhow!("Content not found: {genesis}"))?; + let current = repo + .dag + .get_node(&head)? + .ok_or_else(|| anyhow::anyhow!("Head node not found: {head}"))?; + let mut payload = current.payload().clone(); + let mut policy = policy; + if let Some(previous) = &payload.access_policy { + policy.raise_min_valid_issued_at(previous.min_valid_issued_at()); + } + // Only the policy changed. Keep both the body and its original order. + payload.access_policy = Some(policy); let op = Operation::new(genesis, OperationType::Update(payload), author.to_string()); - - let version_cid = { - let mut repo = self.repo.lock(); - repo.commit_operation(op) - .map_err(|e| anyhow::anyhow!("Failed to commit access policy update: {}", e))? - }; + let version_cid = repo + .commit_operation(op) + .map_err(|e| anyhow::anyhow!("Failed to commit access policy update: {}", e))?; Ok(CommitResult { genesis_cid: genesis_cid.to_string(), @@ -435,110 +437,176 @@ impl ContentRepository for CrslCrdtRepository { .get_operations_with_index(&genesis) .map_err(|e| anyhow::anyhow!("Failed to get operations: {}", e))?; - // The linear history is ordered genesis → head, and `indexed_ops` is - // ordered by timestamp, so position N in one is position N in the - // other. Every lookup below uses that correspondence rather than - // comparing timestamps: two versions written inside the same second - // are indistinguishable by timestamp, and guessing between them - // silently dropped or re-sent versions. - let history = repo - .linear_history(&genesis) - .map_err(|e| anyhow::anyhow!("Failed to get history: {}", e))?; + use crsl_lib::dasl::node::Node; + use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet}; - // Filter by since_version if provided: return only what came after it. - let since_index = if let Some(since) = since_version { - let since_cid = Self::parse_cid(since)?; - - // 1-based, matching `indexed_ops`. An unknown version yields None, - // which sends the full history — the safe direction, since the - // receiver can discard what it already has but cannot invent what - // it never received. - history - .iter() - .position(|cid| *cid == since_cid) - .map(|pos| pos + 1) - } else { - None - }; + let node_cids = repo + .dag + .get_nodes_by_genesis(&genesis) + .map_err(|e| anyhow::anyhow!("Failed to get DAG nodes: {}", e))?; + let mut nodes = HashMap::new(); + for cid in node_cids { + let node = repo + .dag + .get_node(&cid)? + .ok_or_else(|| anyhow::anyhow!("Missing DAG node {cid}"))?; + nodes.insert(cid, node); + } - // Pair each operation with the DAG node it produced. - // - // The receiver recomputes a node's CID from `node_timestamp`, so an - // operation carrying the wrong one is re-derived as a different node — - // or, when two operations carry the same one, collapses onto a single - // CID and the newer version is lost. - // - // `linear_history` is ordered genesis → head and `indexed_ops` is - // ordered by timestamp, so the two line up position by position. That - // is the only reliable correspondence: matching by timestamp proximity - // returned the first node within ±1s, which is the genesis for every - // operation written in the same second. - // A node we cannot read must not be skipped: dropping one shifts every - // pair after it by one, which is worse than the tail simply running - // short. - let node_timestamps = history - .iter() - .map(|cid| { - repo.dag - .get_node(cid) - .ok() - .flatten() - .map(|node| node.timestamp()) - .ok_or_else(|| { - anyhow::anyhow!("Missing DAG node {cid} in the history of {genesis_cid}") - }) - }) - .collect::>>()?; - - // The pairing below is positional, and that only holds while the - // history is a line: `get_operations_with_index` returns every - // operation for this genesis, while `linear_history` picks a single - // child at each branch. Once the DAG forks, the lists differ and an - // index no longer identifies the node an operation produced. - // - // Refuse rather than guess. Handing an operation a timestamp from an - // unrelated node is what corrupted replicas in the first place, and a - // failed sync round is retried; a silently mis-stamped one is not. - if node_timestamps.len() != indexed_ops.len() { - return Err(anyhow::anyhow!( - "Cannot pair operations with their DAG nodes for {genesis_cid}: \ - {} operations over a {}-node linear history. The history is not \ - linear, so no positional pairing is trustworthy.", - indexed_ops.len(), - node_timestamps.len() - )); + // Local operations (including auto-merges) do not store their node + // timestamp. Match the complete node structure with only that field + // removed; operation timestamps and linear-history positions are NOT + // node identities. Imported operations already carry an exact stamp. + let mut by_structure: HashMap> = HashMap::new(); + for (cid, node) in &nodes { + let mut unstamped = node.clone(); + unstamped.timestamp = 0; + by_structure + .entry(unstamped.content_id()?) + .or_default() + .push(*cid); } - let mut operations = Vec::new(); - for (idx, op) in indexed_ops { - // `idx` is 1-based over the full operation list, which is exactly - // the position of this operation's node in the linear history. - // Guarded above: the two lists are the same length, so every - // 1-based operation index addresses its own node. - let node_timestamp = node_timestamps[idx - 1]; - - // Skip operations at or before the since_version index. This runs - // after the lookup above so that skipping never shifts the - // remaining operations onto the wrong nodes. - if let Some(since_idx) = since_index { - if idx <= since_idx { - continue; + let mut templates = Vec::new(); + for (_, op) in indexed_ops { + let node = match &op.kind { + OperationType::Create(payload) => { + Node::new_genesis(payload.clone(), 0, ContentMetadata::default()) } - } + OperationType::Update(payload) | OperationType::Merge(payload) => { + let parent = op + .parents + .first() + .and_then(|cid| nodes.get(cid)) + .ok_or_else(|| anyhow::anyhow!("Missing parent for operation {}", op.id))?; + Node::new_child( + payload.clone(), + op.parents.clone(), + genesis, + 0, + parent.metadata().clone(), + ) + } + // Delete has no payload and the public content API does not + // produce it. Do not guess its historical payload. + OperationType::Delete => { + anyhow::bail!("Cannot pair delete operation {} with its DAG node", op.id) + } + }; + templates.push((op, node)); + } - // Serialize the operation using serde_json for network transfer - let serialized = serde_json::to_vec(&op) - .map_err(|e| anyhow::anyhow!("Failed to serialize operation: {}", e))?; + // Reserve exact imported matches first. Two replicas can independently + // merge identical parents/payloads at different times. Once exchanged, + // the imported stamp disambiguates the remaining local merge node. + let mut claimed = HashSet::new(); + let mut paired = Vec::new(); + let mut local = Vec::new(); + for (op, mut node) in templates { + if let Some(timestamp) = op.node_timestamp { + node.timestamp = timestamp; + let cid = node.content_id()?; + anyhow::ensure!( + nodes.get(&cid) == Some(&node), + "Missing DAG node for operation {}", + op.id + ); + claimed.insert(cid); + paired.push((cid, op)); + } else { + local.push((op, node)); + } + } + for (op, node) in local { + let candidates: Vec<_> = by_structure + .get(&node.content_id()?) + .into_iter() + .flatten() + .filter(|cid| !claimed.contains(*cid)) + .copied() + .collect(); + anyhow::ensure!( + candidates.len() == 1, + "Cannot uniquely pair operation {} with its DAG node: {} candidates", + op.id, + candidates.len() + ); + let cid = candidates[0]; + claimed.insert(cid); + paired.push((cid, op)); + } - operations.push(SerializedOperation { - data: serialized, - genesis_cid: genesis_cid.to_string(), - author: op.author.clone(), - timestamp: op.timestamp, - node_timestamp, - }); + // A receiver holding one branch tip knows its ancestors, not sibling + // branches. Unknown versions safely request the complete DAG. + let mut known = HashSet::new(); + if let Some(since) = since_version { + let since = Self::parse_cid(since)?; + if nodes.contains_key(&since) { + let mut pending = vec![since]; + while let Some(cid) = pending.pop() { + if known.insert(cid) { + let node = nodes + .get(&cid) + .ok_or_else(|| anyhow::anyhow!("Missing ancestor DAG node {cid}"))?; + pending.extend(node.parents().iter().copied()); + } + } + } } + // Emit parents before children even when operation clocks disagree. + // BTreeMap makes the choice between ready siblings deterministic. + let mut pending: BTreeMap>> = BTreeMap::new(); + for (cid, op) in paired { + if !known.contains(&cid) { + pending.entry(cid).or_default().push(op); + } + } + let mut children: HashMap> = HashMap::new(); + let mut remaining_parents = HashMap::new(); + let mut ready_nodes = BTreeSet::new(); + for cid in pending.keys() { + let mut count = 0; + for parent in nodes[cid].parents() { + if !known.contains(parent) { + anyhow::ensure!( + pending.contains_key(parent), + "Missing parent operation for DAG node {parent}" + ); + children.entry(*parent).or_default().push(*cid); + count += 1; + } + } + remaining_parents.insert(*cid, count); + if count == 0 { + ready_nodes.insert(*cid); + } + } + let mut operations = Vec::new(); + while let Some(ready) = ready_nodes.pop_first() { + for op in pending.remove(&ready).expect("ready node is pending") { + operations.push(SerializedOperation { + data: serde_json::to_vec(&op).context("Failed to serialize operation")?, + genesis_cid: genesis_cid.to_string(), + author: op.author.clone(), + timestamp: op.timestamp, + node_timestamp: nodes[&ready].timestamp(), + }); + } + if let Some(children) = children.get(&ready) { + for child in children { + let remaining = remaining_parents.get_mut(child).expect("child is pending"); + *remaining -= 1; + if *remaining == 0 { + ready_nodes.insert(*child); + } + } + } + } + anyhow::ensure!( + pending.is_empty(), + "Cannot export cyclic DAG for {genesis_cid}" + ); Ok(operations) } @@ -651,8 +719,10 @@ impl ContentRepository for CrslCrdtRepository { // 1. Build the Create operation and compute genesis CID via pure math // (Node serialization + SHA-256). No storage is touched. let placeholder = Self::generate_placeholder_cid(data); + let body_updated_at = next_monotonic_timestamp(); let create_payload = ContentPayload { data: data.to_vec(), + body_updated_at, access_policy: None, }; let create_op = Operation::new( @@ -694,6 +764,7 @@ impl ContentRepository for CrslCrdtRepository { let policy = AccessPolicy::new(content_id_vo, identity); let update_payload = ContentPayload { data: data.to_vec(), + body_updated_at, access_policy: Some(policy), }; let update_op = Operation::new( @@ -990,22 +1061,10 @@ mod tests { ); } - /// The position pairing only holds while the history is linear. - /// - /// `get_operations_with_index` returns every operation for the genesis; - /// `linear_history` walks one path and picks a single child at each - /// branch. The moment the DAG forks the two lists differ in length and - /// order, and pairing by index hands an operation a timestamp belonging - /// to some unrelated node — the same "wrong stamp, wrong CID" failure - /// this module was changed to fix, arriving through another door. - /// - /// A node lookup that fails mid-history is worse still: `filter_map` - /// drops it, so every pair after that point shifts by one. - /// - /// We cannot honestly stamp operations we cannot place, so refuse rather - /// than guess. + /// Export all branches using their actual DAG identities, never the + /// positions in the single path returned by linear_history. #[tokio::test] - async fn refuses_to_stamp_operations_it_cannot_place() { + async fn forked_operations_replay_with_their_actual_node_cids() { let (creator, _tmp, genesis_cid) = creator_with_three_versions().await; // Sanity: the linear case still works. @@ -1063,23 +1122,32 @@ mod tests { ) }; - if ops_len == history_len { - // Auto-merge collapsed the fork back into a line; the invariant - // still holds and there is nothing to refuse. - return; - } - - // The lists disagree, so no index pairing is trustworthy. Stamping - // anyway is what corrupts replicas, so the call must fail loudly. - let err = creator.get_operations(&genesis_cid, None).await.expect_err( - "with {ops_len} operations over a {history_len}-node history, \ - pairing by index is a guess and must be refused", - ); - let msg = err.to_string(); - assert!( - msg.contains("history") || msg.contains("operations"), - "the error should say the pairing broke, got: {msg}" + assert!(ops_len > history_len, "fixture must contain a real fork"); + let operations = creator.get_operations(&genesis_cid, None).await.unwrap(); + assert_eq!(operations.len(), ops_len); + let replica_tmp = tempdir().unwrap(); + let replica = CrslCrdtRepository::open(replica_tmp.path()).unwrap(); + assert_eq!( + replica.apply_operations(&operations).await.unwrap(), + ops_len ); + let expected_nodes = { + let repo = creator.repo.lock(); + repo.dag.get_nodes_by_genesis(&genesis).unwrap() + }; + for cid in expected_nodes { + assert_eq!( + replica + .get_version_node_bytes(&genesis_cid, &cid.to_string()) + .await + .unwrap(), + creator + .get_version_node_bytes(&genesis_cid, &cid.to_string()) + .await + .unwrap(), + "every branch node must retain its original CID" + ); + } } /// A replica that committed a WRONG node CID under the old pairing must diff --git a/monas-state-node/src/infrastructure/merge_policy.rs b/monas-state-node/src/infrastructure/merge_policy.rs index 3d90c94..92cefad 100644 --- a/monas-state-node/src/infrastructure/merge_policy.rs +++ b/monas-state-node/src/infrastructure/merge_policy.rs @@ -17,11 +17,12 @@ //! LWW copies *that* whole, and a legitimate write that was never in //! conflict with anything is silently dropped. //! -//! [`MonasMergePolicy`] merges field by field instead: the body is LWW among -//! the heads that actually changed it against their parent (a head that only -//! re-committed its parent's body did not write and must not win a write -//! race), and `min_valid_issued_at` is the max across all heads. Owner and -//! content id are fixed at genesis and simply carried. +//! [`MonasMergePolicy`] merges field by field: the body carries the order of +//! its last explicit write (`body_updated_at`). Policy-only nodes and merge +//! nodes carry that order unchanged, so neither can hide a branch's write or +//! promote an old write using a fresh merge timestamp. Equal write orders use +//! the body bytes as a deterministic tie-break. `min_valid_issued_at` is the +//! max across all heads. Owner and content id are fixed at genesis. //! //! This decides what a merge node *contains*; it does not judge whether a //! head should have been accepted. A write a stale member let through under @@ -39,59 +40,31 @@ use crsl_lib::convergence::policy::{MergePolicy, ResolveInput}; use super::crdt_repository::ContentPayload; -/// Field-wise merge for [`ContentPayload`]: body is LWW among heads that -/// changed it, `min_valid_issued_at` is the max, the rest is carried. +/// Field-wise merge: max body-write register and max invalidation policy. +/// Neither decision depends on the timestamp or parents of a merge node. pub struct MonasMergePolicy; impl MonasMergePolicy { - /// Did this head change the body against its parent(s)? - /// - /// A head with no parent payloads (a genesis, or ancestry the replica - /// has not synced) is taken as a change: there is nothing to say it - /// isn't, and the alternative — dropping it from the race — could lose a - /// real write. - fn changed_body(input: &ResolveInput) -> bool { - input.parent_payloads.is_empty() - || input - .parent_payloads - .iter() - .any(|parent| parent.data != input.payload.data) - } - - /// The head whose body should win: newest among those that wrote. fn body_winner(nodes: &[ResolveInput]) -> &ResolveInput { - let mut writers = nodes.iter().filter(|n| Self::changed_body(n)).peekable(); - if writers.peek().is_some() { - writers - .max_by_key(|n| n.timestamp) - .expect("peeked non-empty") - } else { - // Nobody changed the body: every head carries the same one, so - // any head's copy is right. Newest, for determinism. - nodes - .iter() - .max_by_key(|n| n.timestamp) - .expect("MonasMergePolicy requires at least one candidate node") - } + nodes + .iter() + .max_by(|a, b| { + a.payload + .body_updated_at + .cmp(&b.payload.body_updated_at) + .then_with(|| a.payload.data.cmp(&b.payload.data)) + }) + .expect("MonasMergePolicy requires at least one candidate node") } - /// The policy to carry: the winner's, with the cutoff raised to the max - /// seen on any head. A head without a policy (legacy content) contributes - /// nothing. - fn merged_policy( - nodes: &[ResolveInput], - base: Option<&AccessPolicy>, - ) -> Option { - let base = base.or_else(|| nodes.iter().find_map(|n| n.payload.access_policy.as_ref()))?; - let max_cutoff = nodes + /// Policies share immutable owner/content identity. Select the highest + /// invalidation value independently of the body, retaining its metadata. + fn merged_policy(nodes: &[ResolveInput]) -> Option { + nodes .iter() .filter_map(|n| n.payload.access_policy.as_ref()) - .map(AccessPolicy::min_valid_issued_at) - .max() - .unwrap_or(0); - let mut merged = base.clone(); - merged.raise_min_valid_issued_at(max_cutoff); - Some(merged) + .max_by_key(|p| (p.min_valid_issued_at(), p.updated_at())) + .cloned() } } @@ -100,7 +73,8 @@ impl MergePolicy for MonasMergePolicy { let winner = Self::body_winner(nodes); ContentPayload { data: winner.payload.data.clone(), - access_policy: Self::merged_policy(nodes, winner.payload.access_policy.as_ref()), + body_updated_at: winner.payload.body_updated_at, + access_policy: Self::merged_policy(nodes), } } @@ -133,6 +107,7 @@ mod tests { fn payload(data: &str, cutoff: u64) -> ContentPayload { ContentPayload { data: data.as_bytes().to_vec(), + body_updated_at: 0, access_policy: Some(policy(cutoff)), } } @@ -144,12 +119,12 @@ mod tests { ts: u64, parent: Option, ) -> ResolveInput { - ResolveInput::with_parents( - cid(label), - payload(data, cutoff), - ts, - parent.into_iter().collect(), - ) + let mut value = payload(data, cutoff); + value.body_updated_at = parent + .as_ref() + .filter(|p| p.data == value.data) + .map_or(ts, |p| p.body_updated_at); + ResolveInput::with_parents(cid(label), value, ts, parent.into_iter().collect()) } fn cutoff_of(p: &ContentPayload) -> u64 { @@ -239,4 +214,62 @@ mod tests { assert_eq!(ab, ba); } + + #[test] + fn write_register_is_associative_commutative_and_idempotent() { + let a = head("a", "old", 300, 10, None); + let b = head("b", "latest", 100, 30, None); + let c = head("c", "middle", 200, 20, None); + let expected = MonasMergePolicy.resolve(&[a.clone(), b.clone(), c.clone()]); + assert_eq!(expected.data, b"latest"); + assert_eq!(expected.body_updated_at, 30); + assert_eq!(cutoff_of(&expected), 300); + for [x, y, z] in [ + [a.clone(), b.clone(), c.clone()], + [a.clone(), c.clone(), b.clone()], + [b.clone(), a.clone(), c.clone()], + [b.clone(), c.clone(), a.clone()], + [c.clone(), a.clone(), b.clone()], + [c, b, a], + ] { + let xy = MonasMergePolicy.resolve(&[x.clone(), y.clone()]); + // Deliberately much newer NODE timestamp, preserving BODY order. + let merged = ResolveInput::with_parents( + cid("merge"), + xy.clone(), + 9999, + vec![x.payload, y.payload], + ); + assert_eq!(MonasMergePolicy.resolve(&[merged.clone(), z]), expected); + assert_eq!(MonasMergePolicy.resolve(&[merged.clone(), merged]), xy); + } + } + + #[test] + fn simultaneous_writes_use_body_bytes_not_head_order_or_node_time() { + let a = head("a", "aaa", 100, 10, None); + let mut b = head("b", "zzz", 200, 10, None); + b.timestamp = 1; // it is body_updated_at, not this timestamp, that matters + let ab = MonasMergePolicy.resolve(&[a.clone(), b.clone()]); + let ba = MonasMergePolicy.resolve(&[b, a]); + assert_eq!(ab, ba); + assert_eq!(ab.data, b"zzz"); + assert_eq!(ab.body_updated_at, 10); + assert_eq!(cutoff_of(&ab), 200); + } + + #[test] + fn ancestry_availability_does_not_change_the_body_winner() { + let a = head("a", "new", 100, 30, None); + let b = head("b", "old", 200, 20, None); + let expected = MonasMergePolicy.resolve(&[a.clone(), b.clone()]); + let mut with_parents = a; + with_parents.parent_payloads = vec![with_parents.payload.clone(), b.payload.clone()]; + assert_eq!( + MonasMergePolicy.resolve(&[with_parents.clone(), b.clone()]), + expected + ); + with_parents.parent_payloads.pop(); + assert_eq!(MonasMergePolicy.resolve(&[with_parents, b]), expected); + } } diff --git a/monas-state-node/tests/body_order.rs b/monas-state-node/tests/body_order.rs new file mode 100644 index 0000000..1d972b8 --- /dev/null +++ b/monas-state-node/tests/body_order.rs @@ -0,0 +1,156 @@ +use crsl_lib::{ + convergence::metadata::ContentMetadata, + crdt::operation::{Operation, OperationType}, + dasl::node::Node, +}; +use monas_state_node::{ + domain::{access_policy::AccessPolicy, identity::Identity}, + infrastructure::crdt_repository::{ContentPayload, CrslCrdtRepository}, + port::content_repository::ContentRepository, +}; +use tempfile::tempdir; + +async fn payload(repo: &CrslCrdtRepository, genesis: &str) -> ContentPayload { + let (bytes, _) = repo + .get_latest_node_bytes_with_version(genesis) + .await + .unwrap() + .unwrap(); + Node::::from_bytes(&bytes) + .unwrap() + .payload +} + +#[tokio::test] +async fn a_write_advances_past_an_imported_future_stamp_after_restart() { + let source_dir = tempdir().unwrap(); + let target_dir = tempdir().unwrap(); + let source = CrslCrdtRepository::open(source_dir.path()).unwrap(); + let genesis = source + .create_content(b"base", "owner", None) + .await + .unwrap() + .genesis_cid; + source + .update_content(&genesis, b"zzz future", "owner", None) + .await + .unwrap(); + let mut ops = source.get_operations(&genesis, None).await.unwrap(); + let last = ops.last_mut().unwrap(); + let mut op: Operation = serde_json::from_slice(&last.data).unwrap(); + let OperationType::Update(body) = &mut op.kind else { + panic!("expected a write") + }; + let future = u64::MAX / 2; + body.body_updated_at = future; + last.data = serde_json::to_vec(&op).unwrap(); + { + let target = CrslCrdtRepository::open(target_dir.path()).unwrap(); + target.apply_operations(&ops).await.unwrap(); + assert_eq!(payload(&target, &genesis).await.body_updated_at, future); + } + let target = CrslCrdtRepository::open(target_dir.path()).unwrap(); + target + .update_content(&genesis, b"aaa newer", "owner", None) + .await + .unwrap(); + let latest = payload(&target, &genesis).await; + assert!(latest.body_updated_at > future); + assert_eq!(latest.data, b"aaa newer"); + let third_dir = tempdir().unwrap(); + let third = CrslCrdtRepository::open(third_dir.path()).unwrap(); + third + .apply_operations(&target.get_operations(&genesis, None).await.unwrap()) + .await + .unwrap(); + assert_eq!(payload(&third, &genesis).await, latest); +} + +#[tokio::test] +async fn stale_explicit_policies_cannot_lower_the_observed_boundary() { + let dir = tempdir().unwrap(); + let repo = CrslCrdtRepository::open(dir.path()).unwrap(); + let prepared = repo + .prepare_create_operations( + b"base", + "owner", + Some(Identity::user("owner".into()).unwrap()), + ) + .await + .unwrap(); + repo.apply_operations(&prepared.operations).await.unwrap(); + let genesis = prepared.genesis_cid; + let stale = repo.get_access_policy(&genesis).await.unwrap().unwrap(); + let original = payload(&repo, &genesis).await; + let mut newer = stale.clone(); + newer.raise_min_valid_issued_at(100); + repo.update_access_policy(&genesis, newer, "owner") + .await + .unwrap(); + repo.update_access_policy(&genesis, stale.clone(), "owner") + .await + .unwrap(); + let policy_only = payload(&repo, &genesis).await; + assert_eq!(policy_only.data, original.data); + assert_eq!(policy_only.body_updated_at, original.body_updated_at); + assert_eq!( + policy_only.access_policy.unwrap().min_valid_issued_at(), + 100 + ); + repo.update_content(&genesis, b"new", "owner", Some(stale)) + .await + .unwrap(); + let written = payload(&repo, &genesis).await; + assert_eq!(written.data, b"new"); + assert!(written.body_updated_at > original.body_updated_at); + assert_eq!(written.access_policy.unwrap().min_valid_issued_at(), 100); +} + +#[tokio::test] +async fn create_and_initial_policy_share_body_order_and_missing_order_is_rejected() { + let dir = tempdir().unwrap(); + let repo = CrslCrdtRepository::open(dir.path()).unwrap(); + let prepared = repo + .prepare_create_operations( + b"base", + "owner", + Some(Identity::user("owner".into()).unwrap()), + ) + .await + .unwrap(); + let ops = prepared.operations; + let first: Operation = serde_json::from_slice(&ops[0].data).unwrap(); + let second: Operation = serde_json::from_slice(&ops[1].data).unwrap(); + let OperationType::Create(created) = first.kind else { + panic!("expected create") + }; + let OperationType::Update(policy) = second.kind else { + panic!("expected policy") + }; + assert_eq!(created.body_updated_at, policy.body_updated_at); + let mut legacy_json: serde_json::Value = serde_json::from_slice(&ops[0].data).unwrap(); + legacy_json["kind"]["Create"] + .as_object_mut() + .unwrap() + .remove("body_updated_at"); + let err = + serde_json::from_value::>(legacy_json).unwrap_err(); + assert!(err.to_string().contains("body_updated_at")); + + #[derive(Clone, serde::Serialize, serde::Deserialize)] + struct LegacyPayload { + data: Vec, + access_policy: Option, + } + let legacy = Node::new_genesis( + LegacyPayload { + data: b"base".to_vec(), + access_policy: None, + }, + 1, + ContentMetadata::default(), + ); + assert!( + Node::::from_bytes(&legacy.to_bytes().unwrap()).is_err() + ); +} diff --git a/monas-state-node/tests/fieldwise_convergence.rs b/monas-state-node/tests/fieldwise_convergence.rs new file mode 100644 index 0000000..c966414 --- /dev/null +++ b/monas-state-node/tests/fieldwise_convergence.rs @@ -0,0 +1,188 @@ +use monas_state_node::infrastructure::crdt_repository::CrslCrdtRepository; +use monas_state_node::{AccessPolicy, ContentId, ContentRepository, Identity}; +use tempfile::tempdir; + +#[tokio::test] +async fn newer_write_followed_by_policy_update_must_survive_older_remote_write() { + let ta = tempdir().unwrap(); + let tb = tempdir().unwrap(); + let a = CrslCrdtRepository::open(ta.path()).unwrap(); + let b = CrslCrdtRepository::open(tb.path()).unwrap(); + let genesis = a + .create_content(b"base", "alice", None) + .await + .unwrap() + .genesis_cid; + let policy = AccessPolicy::new( + ContentId::new(genesis.clone()).unwrap(), + Identity::user("alice".into()).unwrap(), + ); + a.update_access_policy(&genesis, policy, "alice") + .await + .unwrap(); + b.apply_operations(&a.get_operations(&genesis, None).await.unwrap()) + .await + .unwrap(); + + b.update_content(&genesis, b"older remote write", "bob", None) + .await + .unwrap(); + tokio::time::sleep(std::time::Duration::from_millis(20)).await; + a.update_content(&genesis, b"newer owner write", "alice", None) + .await + .unwrap(); + let mut policy = a.get_access_policy(&genesis).await.unwrap().unwrap(); + policy.invalidate_tokens(); + a.update_access_policy(&genesis, policy, "alice") + .await + .unwrap(); + + let oa = a.get_operations(&genesis, None).await.unwrap(); + let ob = b.get_operations(&genesis, None).await.unwrap(); + a.apply_operations(&ob).await.unwrap(); + b.apply_operations(&oa).await.unwrap(); + let actual_a = a.get_latest(&genesis).await.unwrap().unwrap(); + let actual_b = b.get_latest(&genesis).await.unwrap().unwrap(); + println!( + "a={}, b={}", + String::from_utf8_lossy(&actual_a), + String::from_utf8_lossy(&actual_b) + ); + assert!( + a.get_access_policy(&genesis) + .await + .unwrap() + .unwrap() + .min_valid_issued_at() + > 0 + ); + assert_eq!( + actual_a, b"newer owner write", + "policy-only head hid this branch's latest real write" + ); + assert_eq!(actual_b, b"newer owner write"); +} + +async fn replicate(from: &CrslCrdtRepository, to: &CrslCrdtRepository, genesis: &str) { + let operations = from.get_operations(genesis, None).await.unwrap(); + assert_eq!( + to.apply_operations(&operations).await.unwrap(), + operations.len() + ); +} + +async fn body(repo: &CrslCrdtRepository, genesis: &str) -> Vec { + repo.get_latest(genesis).await.unwrap().unwrap() +} + +async fn raise_policy(repo: &CrslCrdtRepository, genesis: &str, cutoff: u64) { + let mut policy = repo.get_access_policy(genesis).await.unwrap().unwrap(); + policy.raise_min_valid_issued_at(cutoff); + repo.update_access_policy(genesis, policy, "alice") + .await + .unwrap(); +} + +async fn create(repo: &CrslCrdtRepository) -> String { + let prepared = repo + .prepare_create_operations( + b"base", + "alice", + Some(Identity::user("alice".into()).unwrap()), + ) + .await + .unwrap(); + assert_eq!( + repo.apply_operations(&prepared.operations).await.unwrap(), + prepared.operations.len() + ); + prepared.genesis_cid +} + +#[tokio::test] +async fn two_policy_only_tips_keep_the_latest_body_not_the_latest_revoke() { + let ta = tempdir().unwrap(); + let tb = tempdir().unwrap(); + let a = CrslCrdtRepository::open(ta.path()).unwrap(); + let b = CrslCrdtRepository::open(tb.path()).unwrap(); + let genesis = create(&a).await; + replicate(&a, &b, &genesis).await; + b.update_content(&genesis, b"old", "bob", None) + .await + .unwrap(); + a.update_content(&genesis, b"new", "alice", None) + .await + .unwrap(); + raise_policy(&a, &genesis, 1000).await; + raise_policy(&b, &genesis, 2000).await; + let oa = a.get_operations(&genesis, None).await.unwrap(); + let ob = b.get_operations(&genesis, None).await.unwrap(); + a.apply_operations(&ob).await.unwrap(); + b.apply_operations(&oa).await.unwrap(); + for repo in [&a, &b] { + assert_eq!(body(repo, &genesis).await, b"new"); + assert_eq!( + repo.get_access_policy(&genesis) + .await + .unwrap() + .unwrap() + .min_valid_issued_at(), + 2000 + ); + } +} + +#[tokio::test] +async fn repeated_merges_and_delayed_delivery_keep_original_write_order_after_restart() { + let ta = tempdir().unwrap(); + let tb = tempdir().unwrap(); + let tc = tempdir().unwrap(); + let a = CrslCrdtRepository::open(ta.path()).unwrap(); + let b = CrslCrdtRepository::open(tb.path()).unwrap(); + let c = CrslCrdtRepository::open(tc.path()).unwrap(); + let genesis = create(&a).await; + replicate(&a, &b, &genesis).await; + replicate(&a, &c, &genesis).await; + b.update_content(&genesis, b"Z oldest", "bob", None) + .await + .unwrap(); + c.update_content(&genesis, b"Y delayed", "carol", None) + .await + .unwrap(); + a.update_content(&genesis, b"X newest", "alice", None) + .await + .unwrap(); + let oa = a.get_operations(&genesis, None).await.unwrap(); + let ob = b.get_operations(&genesis, None).await.unwrap(); + a.apply_operations(&ob).await.unwrap(); + b.apply_operations(&oa).await.unwrap(); + assert_eq!(body(&a, &genesis).await, b"X newest"); + assert_eq!(body(&b, &genesis).await, b"X newest"); + // Both replicas mint their own merge of the same concurrent writes. + // Exchange those nodes, then merge the merge nodes before Y arrives. + let ma = a.get_operations(&genesis, None).await.unwrap(); + let mb = b.get_operations(&genesis, None).await.unwrap(); + a.apply_operations(&mb).await.unwrap(); + b.apply_operations(&ma).await.unwrap(); + assert_eq!(body(&a, &genesis).await, b"X newest"); + assert_eq!(body(&b, &genesis).await, b"X newest"); + drop(a); + let a = CrslCrdtRepository::open(ta.path()).unwrap(); + // No policy changes here: an older delayed write must not beat an + // equal-body merge-of-merges, even after its origin is only on disk. + replicate(&c, &a, &genesis).await; + replicate(&c, &b, &genesis).await; + assert_eq!(body(&a, &genesis).await, b"X newest"); + assert_eq!(body(&b, &genesis).await, b"X newest"); + replicate(&a, &c, &genesis).await; + assert_eq!(body(&c, &genesis).await, b"X newest"); + // A subsequent real write must win, not the freshly timestamped merges. + c.update_content(&genesis, b"final explicit write", "carol", None) + .await + .unwrap(); + replicate(&c, &a, &genesis).await; + replicate(&c, &b, &genesis).await; + for repo in [&a, &b, &c] { + assert_eq!(body(repo, &genesis).await, b"final explicit write"); + } +} diff --git a/monas-state-node/tests/integration_test.rs b/monas-state-node/tests/integration_test.rs index 53b8034..b2fd075 100644 --- a/monas-state-node/tests/integration_test.rs +++ b/monas-state-node/tests/integration_test.rs @@ -656,6 +656,7 @@ async fn test_access_control_update_and_verify() { /// のか区別できなくなるため。 #[tokio::test] async fn test_signed_mutation_cannot_be_replayed_after_a_newer_one() { + let request_timestamp = test_timestamp(); let (service, _crdt_repo, _temp_dir) = create_test_service_with_ac().await; service.init_access_control("content-1").await.unwrap(); @@ -681,7 +682,7 @@ async fn test_signed_mutation_cannot_be_replayed_after_a_newer_one() { &update_a, Some(&test_token()), Some(&sig_a), - test_timestamp(), + request_timestamp, ) .await .expect("first update should apply"); @@ -690,7 +691,7 @@ async fn test_signed_mutation_cannot_be_replayed_after_a_newer_one() { &update_b, Some(&test_token()), Some(&sig_b), - test_timestamp(), + request_timestamp, ) .await .expect("second update should apply"); @@ -707,12 +708,14 @@ async fn test_signed_mutation_cannot_be_replayed_after_a_newer_one() { // 攻撃者が捕まえておいた B の署名を再送する。 // (A の再送はドメイン側の単調性チェックが別途弾くので、消費記録が // 効いていることを見るにはこちらを使う) + // Cross a whole-second boundary: replay must preserve the original timestamp. + tokio::time::sleep(std::time::Duration::from_millis(1100)).await; let replay = service .update_access_control( &update_b, Some(&test_token()), Some(&sig_b), - test_timestamp(), + request_timestamp, ) .await; assert!(replay.is_err(), "replaying a consumed signature must fail"); @@ -743,6 +746,7 @@ async fn test_signed_mutation_cannot_be_replayed_after_a_newer_one() { /// 状態巻き戻しが成立する。ID は署名対象メッセージから導くのでこれは効かない。 #[tokio::test] async fn test_signed_mutation_replay_survives_signature_malleability() { + let request_timestamp = test_timestamp(); use p256::ecdsa::Signature; let (service, _crdt_repo, _temp_dir) = create_test_service_with_ac().await; @@ -773,18 +777,20 @@ async fn test_signed_mutation_replay_survives_signature_malleability() { &update, Some(&test_token()), Some(&canonical), - test_timestamp(), + request_timestamp, ) .await .expect("first presentation should apply"); // 同じリクエストを、別バイト列の署名で再送する + // Cross a whole-second boundary: replay must preserve the original timestamp. + tokio::time::sleep(std::time::Duration::from_millis(1100)).await; let replay = service .update_access_control( &update, Some(&test_token()), Some(&malleated), - test_timestamp(), + request_timestamp, ) .await; assert!( @@ -798,6 +804,7 @@ async fn test_signed_mutation_replay_survives_signature_malleability() { /// 通してしまうと正規リクエスト後に発行された Token まで巻き添えで失効する。 #[tokio::test] async fn test_signed_mutation_is_single_use() { + let request_timestamp = test_timestamp(); let (service, _crdt_repo, _temp_dir) = create_test_service_with_ac().await; service.init_access_control("content-1").await.unwrap(); @@ -812,17 +819,19 @@ async fn test_signed_mutation_is_single_use() { &update, Some(&test_token()), Some(&request_signature), - test_timestamp(), + request_timestamp, ) .await .expect("first presentation should apply"); + // Cross a whole-second boundary: replay must preserve the original timestamp. + tokio::time::sleep(std::time::Duration::from_millis(1100)).await; let replay = service .update_access_control( &update, Some(&test_token()), Some(&request_signature), - test_timestamp(), + request_timestamp, ) .await; assert!(replay.is_err(), "the same signature must not apply twice"); diff --git a/monas-state-node/tests/operation_export.rs b/monas-state-node/tests/operation_export.rs new file mode 100644 index 0000000..4258663 --- /dev/null +++ b/monas-state-node/tests/operation_export.rs @@ -0,0 +1,219 @@ +use cid::Cid; +use crsl_lib::{ + convergence::metadata::ContentMetadata, + crdt::operation::{Operation, OperationType}, + dasl::node::Node, +}; +use monas_state_node::{ + infrastructure::crdt_repository::{ContentPayload, CrslCrdtRepository}, + port::content_repository::{ContentRepository, SerializedOperation}, +}; +use tempfile::tempdir; + +// Check the actual DAG CID and bytes, not just payload convergence: a wrong +// timestamp can appear to converge while silently creating a different node. +async fn verify_export( + repo: &CrslCrdtRepository, + genesis: &str, + ops: &[SerializedOperation], +) -> Vec { + let mut cids = Vec::new(); + for wire in ops { + let op: Operation = serde_json::from_slice(&wire.data).unwrap(); + let node = match op.kind { + OperationType::Create(payload) => { + Node::new_genesis(payload, wire.node_timestamp, ContentMetadata::default()) + } + OperationType::Update(payload) | OperationType::Merge(payload) => { + let bytes = repo + .get_version_node_bytes(genesis, &op.parents[0].to_string()) + .await + .unwrap() + .unwrap(); + let parent = Node::::from_bytes(&bytes).unwrap(); + Node::new_child( + payload, + op.parents, + op.genesis, + wire.node_timestamp, + parent.metadata, + ) + } + OperationType::Delete => panic!("fixture does not delete"), + }; + let cid = node.content_id().unwrap().to_string(); + assert_eq!( + repo.get_version_node_bytes(genesis, &cid).await.unwrap(), + Some(node.to_bytes().unwrap()) + ); + cids.push(cid); + } + cids +} + +#[tokio::test] +async fn ambiguous_unstamped_nodes_fail_closed() { + use crsl_lib::{ + graph::storage::{LeveldbNodeStorage, NodeStorage}, + storage::SharedLeveldb, + }; + + let dir = tempdir().unwrap(); + let repo = CrslCrdtRepository::open(dir.path()).unwrap(); + let genesis = repo + .create_content(b"base", "a", None) + .await + .unwrap() + .genesis_cid; + let tip = repo + .update_content(&genesis, b"updated", "a", None) + .await + .unwrap() + .version_cid; + let bytes = repo + .get_version_node_bytes(&genesis, &tip) + .await + .unwrap() + .unwrap(); + drop(repo); + + // Simulate a damaged store with an extra indistinguishable local node. + // Neither operation time nor storage ordering justifies choosing a stamp. + { + let db = SharedLeveldb::open(dir.path().join("crdt_db")).unwrap(); + let store = LeveldbNodeStorage::::new(db); + let mut duplicate = Node::::from_bytes(&bytes).unwrap(); + duplicate.timestamp += 1; + store.put(&duplicate).unwrap(); + } + let repo = CrslCrdtRepository::open(dir.path()).unwrap(); + let err = repo.get_operations(&genesis, None).await.unwrap_err(); + assert!(err.to_string().contains("2 candidates"), "{err}"); +} + +#[tokio::test] +async fn since_a_branch_tip_exports_siblings_and_merge_but_not_ancestors() { + let a_dir = tempdir().unwrap(); + let b_dir = tempdir().unwrap(); + let a = CrslCrdtRepository::open(a_dir.path()).unwrap(); + let b = CrslCrdtRepository::open(b_dir.path()).unwrap(); + let genesis = a + .create_content(b"base", "a", None) + .await + .unwrap() + .genesis_cid; + b.apply_operations(&a.get_operations(&genesis, None).await.unwrap()) + .await + .unwrap(); + let a_tip = a + .update_content(&genesis, b"a", "a", None) + .await + .unwrap() + .version_cid; + let b_tip = b + .update_content(&genesis, b"b", "b", None) + .await + .unwrap() + .version_cid; + a.apply_operations(&b.get_operations(&genesis, None).await.unwrap()) + .await + .unwrap(); + + // Even before a read triggers auto-merge, the sibling must be exported. + let tail = a.get_operations(&genesis, Some(&a_tip)).await.unwrap(); + assert_eq!( + verify_export(&a, &genesis, &tail).await, + vec![b_tip.clone()] + ); + let (_, merge) = a.get_latest_with_version(&genesis).await.unwrap().unwrap(); + let tail = a.get_operations(&genesis, Some(&b_tip)).await.unwrap(); + assert_eq!( + verify_export(&a, &genesis, &tail).await, + vec![a_tip, merge.clone()] + ); + assert_eq!(b.apply_operations(&tail).await.unwrap(), tail.len()); + assert_eq!( + b.get_latest_with_version(&genesis) + .await + .unwrap() + .unwrap() + .1, + merge + ); + assert!(a + .get_operations(&genesis, Some(&merge)) + .await + .unwrap() + .is_empty()); + + let unrelated = a + .create_content(b"unrelated", "a", None) + .await + .unwrap() + .genesis_cid; + let full = a.get_operations(&genesis, None).await.unwrap(); + let fallback = a.get_operations(&genesis, Some(&unrelated)).await.unwrap(); + assert_eq!( + verify_export(&a, &genesis, &fallback).await, + verify_export(&a, &genesis, &full).await + ); + assert_eq!(full.len(), 4); +} + +#[tokio::test] +async fn forked_auto_merges_can_be_exported_and_replicated_repeatedly() { + let a_dir = tempdir().unwrap(); + let b_dir = tempdir().unwrap(); + let a = CrslCrdtRepository::open(a_dir.path()).unwrap(); + let b = CrslCrdtRepository::open(b_dir.path()).unwrap(); + let genesis = a + .create_content(b"base", "a", None) + .await + .unwrap() + .genesis_cid; + let base = a.get_operations(&genesis, None).await.unwrap(); + assert_eq!(b.apply_operations(&base).await.unwrap(), base.len()); + a.update_content(&genesis, b"a", "a", None).await.unwrap(); + b.update_content(&genesis, b"b", "b", None).await.unwrap(); + + // Snapshot both exports before delivery, so both replicas independently + // create an auto-merge with the same parents and payload but distinct CIDs. + for _ in 0..3 { + let a_ops = a.get_operations(&genesis, None).await.unwrap(); + let b_ops = b.get_operations(&genesis, None).await.unwrap(); + let a_cids = verify_export(&a, &genesis, &a_ops).await; + let b_cids = verify_export(&b, &genesis, &b_ops).await; + assert_eq!(b.apply_operations(&a_ops).await.unwrap(), a_ops.len()); + assert_eq!(a.apply_operations(&b_ops).await.unwrap(), b_ops.len()); + for cid in a_cids.iter().chain(&b_cids) { + assert_eq!( + a.get_version_node_bytes(&genesis, cid).await.unwrap(), + b.get_version_node_bytes(&genesis, cid).await.unwrap() + ); + } + assert_eq!(a.get_latest(&genesis).await.unwrap(), Some(b"b".to_vec())); + assert_eq!(b.get_latest(&genesis).await.unwrap(), Some(b"b".to_vec())); + } + + // A fresh third replica must be able to replay the complete forked DAG in + // one export, including merge nodes whose operation timestamp differs. + let ops = a.get_operations(&genesis, None).await.unwrap(); + assert!(ops.iter().any(|wire| { + let op: Operation = serde_json::from_slice(&wire.data).unwrap(); + matches!(op.kind, OperationType::Merge(_)) && op.timestamp != wire.node_timestamp + })); + let c_dir = tempdir().unwrap(); + let c = CrslCrdtRepository::open(c_dir.path()).unwrap(); + assert_eq!(c.apply_operations(&ops).await.unwrap(), ops.len()); + assert_eq!(c.apply_operations(&ops).await.unwrap(), ops.len()); + for cid in verify_export(&a, &genesis, &ops).await { + assert_eq!( + a.get_version_node_bytes(&genesis, &cid).await.unwrap(), + c.get_version_node_bytes(&genesis, &cid).await.unwrap() + ); + } + assert_eq!( + a.get_latest_with_version(&genesis).await.unwrap(), + c.get_latest_with_version(&genesis).await.unwrap() + ); +}