diff --git a/dev/channel-kit.mjs b/dev/channel-kit.mjs index 9e17f4beb..00bdeac0f 100644 --- a/dev/channel-kit.mjs +++ b/dev/channel-kit.mjs @@ -70,16 +70,20 @@ export function validCanvas(event) { typeof event.content === "string" && Buffer.byteLength(event.content) <= 24 * 1024 && Array.isArray(event.tags) && - event.tags.filter((t) => t[0] === "h").length === 1 && + event.tags.filter((t) => t?.[0] === "h").length === 1 && + event.tags.filter((t) => t?.[0] === "expected-revision").length <= 1 && event.tags.every( (t) => Array.isArray(t) && t.length === 2 && (t[0] === "h" ? /^[0-9a-f]{8}-(?:[0-9a-f]{4}-){3}[0-9a-f]{12}$/.test(t[1]) - : t[0] === "client-id" && - typeof t[1] === "string" && - t[1].length <= 128), + : t[0] === "expected-revision" + ? typeof t[1] === "string" && + (t[1] === "none" || /^[0-9a-f]{64}$/.test(t[1])) + : t[0] === "client-id" && + typeof t[1] === "string" && + t[1].length <= 128), ) ); } diff --git a/dev/channel-kit.test.mjs b/dev/channel-kit.test.mjs index 4d2f6da04..1f5f75e19 100644 --- a/dev/channel-kit.test.mjs +++ b/dev/channel-kit.test.mjs @@ -148,7 +148,23 @@ it("admits only bounded Canvas writes with one exact channel and no notification tags: [["h", "11111111-1111-4111-8111-111111111111"]], }; expect(validCanvas(canvas)).toBe(true); + for (const revision of ["none", "a".repeat(64)]) + expect( + validCanvas({ + ...canvas, + tags: [...canvas.tags, ["expected-revision", revision]], + }), + ).toBe(true); for (const tags of [ + [...canvas.tags, ["expected-revision", "bad"]], + [...canvas.tags, ["expected-revision", "A".repeat(64)]], + [...canvas.tags, ["expected-revision", "none", "extra"]], + [ + ...canvas.tags, + ["expected-revision", "none"], + ["expected-revision", "none"], + ], + [...canvas.tags, null], [], [...canvas.tags, ...canvas.tags], [...canvas.tags, ["p", owner.pubkey]], diff --git a/dev/relay-broker-live.test.mjs b/dev/relay-broker-live.test.mjs index 6c05d6d50..44b92460f 100644 --- a/dev/relay-broker-live.test.mjs +++ b/dev/relay-broker-live.test.mjs @@ -1377,11 +1377,23 @@ test.each([ ); test.each([ - ["conflict: artifact head changed", "failed", "conflict: the relay state"], - ["error: internal server error", "unknown", "could not be confirmed"], + [ + 45010, + "conflict: artifact head changed", + "failed", + "conflict: the relay state", + ], + [45010, "error: internal server error", "unknown", "could not be confirmed"], + [ + 40100, + "conflict: canvas head changed", + "failed", + "conflict: the relay state", + ], + [40100, "error: internal server error", "unknown", "could not be confirmed"], ])( - "artifact refusal %s reaches broker/outbox as %s / %s", - async (reason, delivery, error) => { + "kind %s refusal %s reaches broker/outbox as %s / %s", + async (kind, reason, delivery, error) => { const h = await harness(); let traffic, owner; try { @@ -1394,9 +1406,12 @@ test.each([ save() {}, }); const id = owner.outbox.send({ - kind: 45010, + kind, content: "", - tags: [["h", "00000000-0000-4000-8000-000000000001"]], + tags: [ + ["h", "00000000-0000-4000-8000-000000000001"], + ["expected-revision", "none"], + ], }); await until(() => h.publications.length === 1); await h.sockets[0].receive(["OK", id, false, reason]); diff --git a/docs/plugin-architecture.md b/docs/plugin-architecture.md index 437edabe7..cc77fe6e4 100644 --- a/docs/plugin-architecture.md +++ b/docs/plugin-architecture.md @@ -266,10 +266,26 @@ Task actions pause while saving; the new-item input stays editable and retains its text when the save finishes. If the loaded Canvas was written in the current second, a single cancellable wait respects its timestamp ordering; there is no background retry loop. Failures and recovered drafts expose Retry rather than silently publishing -on reopen. Save uses the existing session Canvas/outbox contract, including its 24 KiB limit, fresh membership check, -optimistic head comparison and exact confirmation. This is **not atomic concurrency -control**; simultaneous saves can overwrite edits. Detected conflicts retain the -local draft and require reviewing the saved Canvas. Refresh confirms before discarding +on reopen. Save uses the existing session Canvas/outbox contract, including its +24 KiB limit, fresh membership check, writer-backed Canvas head/editor-confirmation +reads and exact signed-event recovery. Canvas reads default to strong consistency +for editor/Todos bases and setup preconditions. Setup delivery and its separate +exact-ID confirmation both use writer-backed reads; unknown seed outcomes are +checked without automatically replaying the seed. Template copies explicitly opt +out and remain replica-eligible; Channel Settings no longer reads a Canvas preview. +Editor/Todos saves carry `expected-revision=` (or `none` when absent); template seeds carry `none`. On relays supporting +Canvas compare-and-swap, stale preconditions are refused atomically before mutation. +A proven conflict keeps the local draft and dismisses only that rejected outbox +operation so a reviewed save can proceed. Unknown outcomes remain in Outbox and +block replacement; exact signed retries keep their original precondition. The +post-write different-head check remains conservative and also retains the draft. + +**Compatibility gate:** deploy with a relay supporting Canvas revision preconditions +([block/buzz#6780](https://github.com/block/buzz/pull/6780)) for atomic protection. +An older relay may ignore the tag; strong reads and client head comparison alone +cannot prevent concurrent overwrite. This client change does not upgrade the relay +or add Canvas history/restore. Detected conflicts require reviewing the saved Canvas. Refresh confirms before discarding edits. Local recovery drafts are partitioned by community/viewer/channel; if browser storage is unavailable they survive only while the editor stays open. Save never promotes local recovery storage to shared state. Already accepted outbox operations diff --git a/docs/relay-queries.md b/docs/relay-queries.md index 575f4d67e..5586aed25 100644 --- a/docs/relay-queries.md +++ b/docs/relay-queries.md @@ -96,8 +96,10 @@ just-added member does not fail the next operation. Template setup also confirms exact Canvas/member events and selected Canvas heads against the writer without replaying accepted commands. Agent deletion discovers member channels and confirms each removal with writer-backed rosters; unreadable -rosters still fail closed. Ordinary Canvas browsing and standalone Canvas/recipe -save confirmation routing are unchanged by this policy. +rosters still fail closed. Standalone recipe saves use writer-backed exact-ID +confirmation; recipe head and catalog reads remain replica-eligible. For +writer-backed Canvas editor/Todos reads and replica-eligible template copies, see +the [Canvas/outbox contract](plugin-architecture.md#optional-canvas-todos). ## Community emoji diff --git a/src-tauri/src/relay.rs b/src-tauri/src/relay.rs index bcdc9ac5c..69371021c 100644 --- a/src-tauri/src/relay.rs +++ b/src-tauri/src/relay.rs @@ -215,6 +215,10 @@ fn validate_event(community: &str, event: &EventTemplate) -> Result<()> { if !channel_writes::creation(event) { return Err("Agent enrollment or channel operation unavailable or invalid".into()); } + } else if event.kind == 40100 { + if !valid_canvas(event) { + return Err("Malformed Canvas save".into()); + } } else if event.kind == 28936 { // A NIP-43 leave request revokes the signer's own membership: empty // content and exactly the NIP-70 protected tag, nothing else. @@ -223,25 +227,47 @@ fn validate_event(community: &str, event: &EventTemplate) -> Result<()> { } } else if !matches!( event.kind, - 0 | 7 - | 9 - | 1984 - | 9000 - | 9001 - | 20001 - | 30030 - | 30177 - | 30315 - | 40003 - | 40100 - | 42000 - | 45010 + 0 | 7 | 9 | 1984 | 9000 | 9001 | 20001 | 30030 | 30177 | 30315 | 40003 | 42000 | 45010 ) { return Err("This event is not supported by the packaged relay connection".into()); } Ok(()) } +// Match the broker's purpose-bound Canvas admission, including legacy untagged retries. +fn valid_canvas(event: &EventTemplate) -> bool { + event.content.len() <= 24 * 1024 + && event + .tags + .iter() + .filter(|tag| tag.first().map(String::as_str) == Some("h")) + .count() + == 1 + && event + .tags + .iter() + .filter(|tag| tag.first().map(String::as_str) == Some("expected-revision")) + .count() + <= 1 + && event.tags.iter().all(|tag| { + if tag.len() != 2 { + return false; + } + match tag[0].as_str() { + "h" => channel_writes::uuid(&tag[1]), + "client-id" => tag[1].encode_utf16().count() <= 128, + "expected-revision" => { + tag[1] == "none" + || (tag[1].len() == 64 + && tag[1] + .bytes() + .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))) + } + _ => false, + } + }) +} + /** Keep the shared kind-5 writer aligned with the broker's channel-local deletion shape. */ fn valid_message_deletion(event: &EventTemplate) -> bool { if event.kind != 5 diff --git a/src-tauri/src/relay/tests.rs b/src-tauri/src/relay/tests.rs index 735c68575..f0961f503 100644 --- a/src-tauri/src/relay/tests.rs +++ b/src-tauri/src/relay/tests.rs @@ -1800,3 +1800,56 @@ async fn preference_batches_reject_invalid_ciphertext_after_signature_verificati } } } + +#[test] +fn canvas_signing_bounds_revision_preconditions_and_allows_exact_legacy_retries() { + let channel = vec![ + "h".to_string(), + "11111111-1111-4111-8111-111111111111".to_string(), + ]; + let event = |tags| EventTemplate { + kind: 40100, + content: "# Plan".into(), + created_at: 100, + tags, + }; + assert!(validate_event("https://relay.test", &event(vec![channel.clone()])).is_ok()); + for revision in ["none".to_string(), "a".repeat(64)] { + assert!(validate_event( + "https://relay.test", + &event(vec![ + channel.clone(), + vec!["expected-revision".into(), revision] + ]) + ) + .is_ok()); + } + for tags in [ + vec![], + vec![channel.clone(), channel.clone()], + vec![channel.clone(), vec!["p".into(), "a".repeat(64)]], + vec![ + channel.clone(), + vec!["expected-revision".into(), "bad".into()], + ], + vec![ + channel.clone(), + vec!["expected-revision".into(), "A".repeat(64)], + ], + vec![ + channel.clone(), + vec!["expected-revision".into(), "none".into(), "extra".into()], + ], + vec![ + channel.clone(), + vec!["expected-revision".into(), "none".into()], + vec!["expected-revision".into(), "none".into()], + ], + vec![channel.clone(), vec![]], + ] { + assert!(validate_event("https://relay.test", &event(tags)).is_err()); + } + let mut too_large = event(vec![channel]); + too_large.content = "é".repeat(13 * 1024); + assert!(validate_event("https://relay.test", &too_large).is_err()); +} diff --git a/src/bundled/channel-templates/TemplateSettings.tsx b/src/bundled/channel-templates/TemplateSettings.tsx index d1e4c40d6..e5df3fa2c 100644 --- a/src/bundled/channel-templates/TemplateSettings.tsx +++ b/src/bundled/channel-templates/TemplateSettings.tsx @@ -98,7 +98,7 @@ export function SaveAsTemplate({ setBusy(true); setError(""); try { - const canvas = await session.canvas.read(channel.id); + const canvas = await session.canvas.read(channel.id, { strong: false }); if (!active()) return; setCopyWarning( catalog.agentsComplete diff --git a/src/bundled/channel-templates/agent-selection.test.tsx b/src/bundled/channel-templates/agent-selection.test.tsx index 5e1de2f56..f1339c588 100644 --- a/src/bundled/channel-templates/agent-selection.test.tsx +++ b/src/bundled/channel-templates/agent-selection.test.tsx @@ -515,12 +515,16 @@ it("does not consume a legacy group default while its required identity is still it("copies a complete managed lineup without an unused legacy warning", async () => { const test = harness(), user = userEvent.setup(); + const read = vi.fn(test.owner.session.canvas.read); try { await test.owner.session.agentChoices.refresh(); await test.owner.session.archives.ensure(); render( { + const operation: OutgoingEvent = { + event: signed(viewer, { kind, content: "Draft", tags: [] }), + delivery, + error, + }; + expect(isDefinitiveCanvasConflict(operation)).toBe(expected); + }, +); + +it("does not classify missing delivery evidence as a refusal", () => { + expect(isDefinitiveCanvasConflict(undefined)).toBe(false); +}); diff --git a/src/features/channel-templates/canvas-conflict.ts b/src/features/channel-templates/canvas-conflict.ts new file mode 100644 index 000000000..3b42ea904 --- /dev/null +++ b/src/features/channel-templates/canvas-conflict.ts @@ -0,0 +1,12 @@ +import type { OutgoingEvent } from "../relay/outbox"; + +/** A refused retry cannot prove an earlier uncertain publication never landed. */ +export function isDefinitiveCanvasConflict( + operation: OutgoingEvent | undefined, +): boolean { + return ( + operation?.event.kind === 40100 && + operation.delivery === "failed" && + operation.error?.startsWith("conflict:") === true + ); +} diff --git a/src/features/channel-templates/capability.test.ts b/src/features/channel-templates/capability.test.ts index 247262a65..f67800720 100644 --- a/src/features/channel-templates/capability.test.ts +++ b/src/features/channel-templates/capability.test.ts @@ -5,7 +5,12 @@ import { coordinate, KIT_TAG, type KitRecord } from "./model"; import { keypair, metadata, roster, signed } from "../relay/testing"; import { matchesEvent } from "../relay/projection"; import type { ReadFilter, RelayEvent } from "../relay/events"; -import type { Outbox, OutgoingEvent } from "../relay/outbox"; +import { + createOutbox, + PublishRejected, + type Outbox, + type OutgoingEvent, +} from "../relay/outbox"; import type { RelayWriter } from "../relay/transport"; import type { RelayReader } from "../relay/reader"; @@ -204,7 +209,10 @@ it("keeps drafts safe from stale Canvas, unresolved writes and access loss", asy expect(saved.content).toBe("New"); expect(f.outbox.send).toHaveBeenCalledWith({ kind: 40100, - tags: [["h", channel]], + tags: [ + ["h", channel], + ["expected-revision", "none"], + ], content: "New", }); }); @@ -334,3 +342,268 @@ it.each(["recipe", "Canvas"])( } }, ); + +it("keeps display-only Canvas reads replica-eligible without weakening editor reads", async () => { + const f = fixture(); + await f.canvas.read(channel, { strong: false }); + expect(f.reader.read).toHaveBeenLastCalledWith( + [{ kinds: [40100], "#h": [channel], limit: 1 }], + { signal: f.controller.signal, fresh: true }, + ); + await f.canvas.read(channel); + expect(f.reader.read).toHaveBeenLastCalledWith( + [{ kinds: [40100], "#h": [channel], limit: 1, consistency: "strong" }], + { signal: f.controller.signal, fresh: true }, + ); +}); + +it.each([false, true])( + "carries a strong-read revision precondition for existing=%s", + async (existing) => { + const f = fixture(); + const before = signed(viewer, { + kind: 40100, + tags: [["h", channel]], + content: "Before", + created_at: 100, + }); + if (existing) f.events.push(before); + const base = await f.canvas.read(channel); + const saved = await f.canvas.save(channel, "After", base?.id); + expect(saved.tags).toContainEqual([ + "expected-revision", + existing ? before.id : "none", + ]); + for (const [filters] of vi.mocked(f.reader.read).mock.calls) + expect(filters.every((filter) => filter.consistency === "strong")).toBe( + true, + ); + }, +); + +it.each(["failed", "unknown"] as const)( + "cleans up only proven CAS refusal, not %s uncertainty", + async (delivery) => { + const f = fixture(); + vi.mocked(f.outbox.send).mockImplementation((value) => { + const event = signed(viewer, { + ...value, + created_at: Math.floor(Date.now() / 1000), + }); + f.pending.push({ + event, + delivery, + error: "conflict: the relay state changed", + }); + return event.id; + }); + f.delivered.mockRejectedValue( + new Error("conflict: the relay state changed"), + ); + const save = f.canvas.save(channel, "Draft", undefined); + await expect(save).rejects.toThrow( + delivery === "failed" ? /Canvas changed.*draft is kept/ : /conflict:/, + ); + if (delivery === "failed") { + expect(f.outbox.dismiss).toHaveBeenCalledWith(f.pending[0]?.event.id); + expect(f.delivered).not.toHaveBeenCalled(); + } else expect(f.outbox.dismiss).not.toHaveBeenCalled(); + }, +); + +it.each([false, true])( + "an intervening relay head survives a stale save (existing=%s), then a reviewed save succeeds", + async (existing) => { + const f = fixture(); + let head: RelayEvent | undefined = existing + ? signed(viewer, { + kind: 40100, + tags: [["h", channel]], + content: "Before", + created_at: 100, + }) + : undefined; + const before = head; + const winner = signed(viewer, { + kind: 40100, + tags: [["h", channel]], + content: "Other editor", + created_at: 101, + }); + let race = true; + const filtersSeen: ReadFilter[][] = []; + const reader: RelayReader = { + read: vi.fn(async (filters: readonly ReadFilter[]) => { + filtersSeen.push([...filters]); + return head && + filters.some((filter) => matchesEvent(head as RelayEvent, filter)) + ? [head] + : []; + }), + }; + const publish = vi.fn(async (event: RelayEvent) => { + if (race) { + head = winner; + race = false; + } + if ( + event.tags.find(([name]) => name === "expected-revision")?.[1] !== + (head?.id ?? "none") + ) + throw new PublishRejected("conflict: the relay state changed"); + head = event; + }); + const owner = createOutbox( + viewer.pubkey, + { + kinds: [40100], + sign: async (event) => signed(viewer, event), + publish, + }, + { load: () => [], save() {} }, + ); + const kit = createChannelKit({ + host: undefined, + reader, + outbox: owner.outbox, + local: owner.local, + ready: owner.outbox.ready(), + viewer: viewer.pubkey, + community, + signal: f.controller.signal, + canWrite: () => true, + delivered: async (id) => { + await vi.waitFor(() => + expect( + owner.outbox.snapshot().find((item) => item.event.id === id) + ?.delivery, + ).not.toBe("sending"), + ); + const item = owner.outbox + .snapshot() + .find((item) => item.event.id === id); + if (item?.delivery === "failed" || item?.delivery === "unknown") + throw new Error(item.error); + }, + }); + try { + await expect( + kit.canvas.save(channel, "Stale draft", before?.id), + ).rejects.toThrow(/Canvas changed.*draft is kept/); + expect(head).toBe(winner); + expect(owner.outbox.snapshot()).toEqual([]); + const saved = await kit.canvas.save(channel, "Reviewed draft", winner.id); + expect(saved.content).toBe("Reviewed draft"); + expect(head?.id).toBe(saved.id); + expect(publish).toHaveBeenCalledTimes(2); + expect( + filtersSeen.flat().every((filter) => filter.consistency === "strong"), + ).toBe(true); + } finally { + owner.dispose(); + f.controller.abort(); + } + }, +); + +it("retains an uncertain Canvas after lost ACK and CAS-refused exact retry, blocking replacement", async () => { + const f = fixture(); + const before = signed(viewer, { + kind: 40100, + tags: [["h", channel]], + content: "Before", + created_at: 100, + }); + const winner = signed(viewer, { + kind: 40100, + tags: [["h", channel]], + content: "Other editor", + created_at: 101, + }); + let head = before; + const reader: RelayReader = { + // No exact-ID confirmation is available for the uncertain publication. + read: async (filters) => + filters.some((filter) => !filter.ids) ? [head] : [], + }; + const sign = vi.fn(async (event: Parameters[0]) => + signed(viewer, event), + ); + const publish = vi + .fn(async (_event: RelayEvent) => {}) + .mockRejectedValueOnce(new Error("Lost ACK")) + .mockRejectedValueOnce( + new PublishRejected("conflict: the relay state changed"), + ); + const storage = { load: () => [], save: vi.fn() }; + const owner = createOutbox( + viewer.pubkey, + { kinds: [40100], sign, publish }, + storage, + ); + const dismiss = vi.fn(owner.outbox.dismiss); + const delivered = async (id: string) => { + await vi.waitFor(() => + expect( + owner.outbox.snapshot().find((item) => item.event.id === id)?.delivery, + ).toBe("unknown"), + ); + throw new Error( + owner.outbox.snapshot().find((item) => item.event.id === id)?.error, + ); + }; + const kit = createChannelKit({ + host: undefined, + reader, + outbox: { ...owner.outbox, dismiss }, + local: owner.local, + ready: owner.outbox.ready(), + viewer: viewer.pubkey, + community, + signal: f.controller.signal, + canWrite: () => true, + delivered, + }); + try { + await expect( + kit.canvas.save(channel, "My draft", before.id), + ).rejects.toThrow("Lost ACK"); + const original = publish.mock.calls[0]?.[0]; + expect(original).toBeDefined(); + if (!original) throw new Error("Missing initial publication"); + expect(original.tags).toContainEqual(["expected-revision", before.id]); + expect(owner.outbox.snapshot()[0]).toMatchObject({ + signed: original, + delivery: "unknown", + }); + head = winner; + owner.outbox.retry(original.id); + await vi.waitFor(() => + expect(owner.outbox.snapshot()[0]).toMatchObject({ + delivery: "unknown", + error: "Retry blocked: conflict: the relay state changed", + }), + ); + // Canvas cleanup must respect the outbox's earlier uncertain dispatch, + // even though this later attempt was definitively refused. + await expect(kit.confirm(original.id)).rejects.toThrow( + "Retry blocked: conflict:", + ); + await expect( + kit.canvas.save(channel, "Reviewed draft", winner.id), + ).rejects.toThrow(/unresolved Canvas save/); + expect(dismiss).not.toHaveBeenCalled(); + expect(sign).toHaveBeenCalledTimes(1); + expect(publish.mock.calls.map(([event]) => event)).toEqual([ + original, + original, + ]); + expect(storage.save).toHaveBeenLastCalledWith([ + expect.objectContaining({ signed: original, delivery: "unknown" }), + ]); + expect(head).toBe(winner); + } finally { + owner.dispose(); + f.controller.abort(); + } +}); diff --git a/src/features/channel-templates/capability.ts b/src/features/channel-templates/capability.ts index 753618e4f..705116dee 100644 --- a/src/features/channel-templates/capability.ts +++ b/src/features/channel-templates/capability.ts @@ -2,6 +2,7 @@ import type { RelayReader } from "../relay/reader"; import type { RelayEvent, ReadFilter } from "../relay/events"; import type { Outbox, LocalEvents } from "../relay/outbox"; import type { ChannelKitHost } from "./host"; +import { isDefinitiveCanvasConflict } from "./canvas-conflict"; import { CANVAS_BYTES, coordinate, @@ -12,6 +13,9 @@ import { type KitValue, } from "./model"; +const canvasConflict = + "Canvas changed since you opened it. Your draft is kept; load the current document before replacing it."; + export const selectedHead = (events: readonly RelayEvent[]) => [...events].sort( (a, b) => b.created_at - a.created_at || a.id.localeCompare(b.id), @@ -130,18 +134,38 @@ export function createChannelKit({ return pending; } async function confirm(id: string) { - if ((await fresh([{ ids: [id], limit: 1 }])).some((e) => e.id === id)) - return; - const saved = local?.snapshot().find((e) => e.event.id === id); - if (!saved || Date.now() / 1000 - saved.event.created_at >= 15 * 60) - throw new Error( - "This operation could not be confirmed and is too old to replay. Keep your draft and inspect the current result before making a new save.", - ); - await delivered(id); - if (!(await fresh([{ ids: [id], limit: 1 }])).some((e) => e.id === id)) - throw new Error( - "Save is awaiting exact relay confirmation; your draft is kept", - ); + try { + if ( + (await fresh([{ ids: [id], limit: 1, consistency: "strong" }])).some( + (e) => e.id === id, + ) + ) + return; + const saved = local?.snapshot().find((e) => e.event.id === id); + if (!saved || Date.now() / 1000 - saved.event.created_at >= 15 * 60) + throw new Error( + "This operation could not be confirmed and is too old to replay. Keep your draft and inspect the current result before making a new save.", + ); + // A proven Canvas CAS refusal cannot become valid by replaying these bytes. + if (isDefinitiveCanvasConflict(saved)) throw new Error(saved.error); + await delivered(id); + if ( + !(await fresh([{ ids: [id], limit: 1, consistency: "strong" }])).some( + (e) => e.id === id, + ) + ) + throw new Error( + "Save is awaiting exact relay confirmation; your draft is kept", + ); + } catch (error) { + const failed = local?.snapshot().find((item) => item.event.id === id); + if (outbox && isDefinitiveCanvasConflict(failed)) { + // Refused before mutation. Release the pending gate, not uncertain saves. + await outbox.dismiss(id); + throw new Error(canvasConflict); + } + throw error; + } } const capability = Object.freeze({ available: !!host && !!outbox?.supports(30078), @@ -227,13 +251,18 @@ export function createChannelKit({ }); const canvas = Object.freeze({ available: !!outbox?.supports(40100), - async read(channel: string, options: Pick = {}) { + async read(channel: string, { strong = true }: { strong?: boolean } = {}) { signal.throwIfAborted(); if (!canWrite(channel)) throw new Error("Canvas is unavailable after channel access changed"); return selectedHead( await fresh([ - { kinds: [40100], "#h": [channel], limit: 1, ...options }, + { + kinds: [40100], + "#h": [channel], + limit: 1, + ...(strong ? { consistency: "strong" as const } : {}), + }, ]), ); }, @@ -246,10 +275,7 @@ export function createChannelKit({ await ready; signal.throwIfAborted(); const head = await canvas.read(channel); - if (head?.id !== expected) - throw new Error( - "Canvas changed since you opened it. Your draft is kept; load the current document before replacing it.", - ); + if (head?.id !== expected) throw new Error(canvasConflict); const pending = local ?.snapshot() .find( @@ -265,7 +291,14 @@ export function createChannelKit({ if (head && head.created_at >= Math.floor(Date.now() / 1000)) throw new Error("Please wait a second before saving Canvas again"); if (!canWrite(channel)) throw new Error("Channel access changed"); - const id = outbox.send({ kind: 40100, content, tags: [["h", channel]] }); + const id = outbox.send({ + kind: 40100, + content, + tags: [ + ["h", channel], + ["expected-revision", expected ?? "none"], + ], + }); await confirm(id); const selected = await canvas.read(channel); if (selected?.id !== id) diff --git a/src/features/channel-templates/setup.test.ts b/src/features/channel-templates/setup.test.ts index 652fdc256..0decfb816 100644 --- a/src/features/channel-templates/setup.test.ts +++ b/src/features/channel-templates/setup.test.ts @@ -2,7 +2,10 @@ import { afterEach, beforeEach, assert, expect, it, vi } from "vitest"; import { createChannelSetup, type ChannelCreationInput } from "./setup"; import type { Outbox, OutgoingEvent } from "../relay/outbox"; -import { keypair, signed } from "../relay/testing"; +import { keypair, metadata, roster, signed } from "../relay/testing"; +import { createRelaySession } from "../relay/session"; +import { matchesEvent } from "../relay/projection"; +import type { ReadFilter, RelayEvent } from "../relay/events"; const viewer = keypair(); const agent = keypair().pubkey; @@ -104,6 +107,100 @@ function hold() { return { promise, release }; } +it.each([true, false])( + "confirms a setup seed against the writer without replay (writer visible: %s)", + async (visible) => { + const relay = keypair(); + const events: RelayEvent[] = []; + const published: RelayEvent[] = []; + const query = vi.fn(async (filters: readonly ReadFilter[]) => + events.filter((event) => + filters.some( + (filter) => + matchesEvent(event, filter) && + (event.kind !== 40100 || + (visible && filter.consistency === "strong")), + ), + ), + ); + const owner = createRelaySession( + { + viewer: viewer.pubkey, + relayAuthor: relay.pubkey, + scope: "https://setup.example.test", + media: () => undefined, + query, + channelKit: { prepare: async () => "", decode: async () => [] }, + writer: { + kinds: [9007, 40100], + sign: async (template) => signed(viewer, template), + publish: async (event) => { + published.push(event); + events.push(event); + if (event.kind === 9007) { + const id = event.tags.find(([tag]) => tag === "h")?.[1]; + assert.exists(id); + events.push( + metadata(relay, id, "Daily"), + roster(relay, id, [viewer.pubkey]), + ); + } + }, + }, + }, + { outboxStorage: { load: async () => [], save: async () => {} } }, + ); + try { + const id = await owner.session.channelCreation.create({ + name: "Daily", + visibility: "open", + setup: { canvas: "# Plan", agents: [], groupId: "", templateId: "" }, + }); + const receipt = `buzz-channel-setup.v2:https://setup.example.test:${viewer.pubkey}:${id}`; + if (visible) { + await vi.waitFor(() => + expect(localStorage.getItem(receipt)).toBeNull(), + ); + expect(owner.session.channelCreation.notices()).toEqual([]); + expect(owner.session.outbox?.snapshot()).toEqual([]); + } else { + await vi.waitFor(() => + expect(owner.session.channelCreation.notices()).toEqual([ + expect.objectContaining({ + id, + error: expect.stringContaining( + "awaiting exact relay confirmation", + ), + }), + ]), + ); + expect( + JSON.parse(localStorage.getItem(receipt) ?? "null"), + ).toMatchObject({ + canvasDone: false, + }); + } + // Receipt retirement / completion failure is the barrier: confirmation + // must not republish even when the writer cannot prove the accepted seed. + expect(published.map((event) => event.kind)).toEqual([9007, 40100]); + const seed = published[1]; + assert.exists(seed); + const exact = query.mock.calls.flatMap(([filters]) => + filters.filter((filter) => filter.ids?.includes(seed.id)), + ); + // Outbox delivery may also observe by ID; this pins setup's separate, + // fresh confirmation read rather than incidental background read counts. + expect(exact).toContainEqual({ + ids: [seed.id], + limit: 1, + consistency: "strong", + }); + } finally { + owner.dispose(); + } + }, +); + it("opens after confirmed creation, before Canvas completes; places independently and never sends a message", async () => { const f = fixture(); const gate = hold(); @@ -119,6 +216,10 @@ it("opens after confirmed creation, before Canvas completes; places independentl await reached.promise; expect(f.opts.place).toHaveBeenCalledWith(id, "work", undefined); expect(f.events.map((e) => e.event.kind)).toEqual([9007, 40100]); + expect(f.events[1]?.event.tags).toEqual([ + ["h", id], + ["expected-revision", "none"], + ]); expect(receipts()).toHaveLength(1); gate.release(); await run.completion; @@ -202,6 +303,48 @@ it("Canvas failure does not lose group placement or turn admission into failure" expect(receipts()).toHaveLength(1); }); +it.each(["failed", "unknown"] as const)( + "keeps seed recovery observe-only and dismisses only a proven conflict (%s)", + async (delivery) => { + const f = fixture(); + f.opts.confirm.mockImplementation(async (id) => { + const index = f.events.findIndex((item) => item.event.id === id); + assert.exists(f.events[index]); + f.events[index] = { + ...f.events[index], + delivery, + error: "conflict: the relay state changed", + }; + f.setCanvas("other-editor"); + throw new Error("conflict: the relay state changed"); + }); + const run = f.setup.run(input, viewer.pubkey); + const failure = expect(run.completion).rejects.toThrow( + delivery === "failed" + ? /Canvas changed.*no agents were added/ + : /conflict:/, + ); + await run.admission; + await failure; + expect(f.opts.canvasHead).toHaveBeenCalledTimes(1); + expect(await f.opts.canvasHead()).toBe("other-editor"); + expect(f.outbox.retry).not.toHaveBeenCalled(); + expect( + vi.mocked(f.outbox.send).mock.calls.map(([event]) => event.kind), + ).toEqual([9007, 40100]); + expect(f.events.map((item) => item.event.kind)).toEqual( + delivery === "failed" ? [9007] : [9007, 40100], + ); + expect(f.outbox.dismiss).toHaveBeenCalledTimes( + delivery === "failed" ? 1 : 0, + ); + const saved = JSON.parse(localStorage.getItem(firstReceipt()) ?? "null"); + expect(saved.created).toBe(true); + expect(saved.canvasDone).toBe(false); + expect(saved.added).toEqual([]); + }, +); + it("group failure does not skip Canvas and members; retains the incomplete destination", async () => { const f = fixture(); f.opts.place.mockRejectedValue(new Error("Group unavailable")); diff --git a/src/features/channel-templates/setup.ts b/src/features/channel-templates/setup.ts index 5deda7a8b..b22ecfc2e 100644 --- a/src/features/channel-templates/setup.ts +++ b/src/features/channel-templates/setup.ts @@ -1,5 +1,6 @@ import type { Outbox, LocalEvents } from "../relay/outbox"; import type { KitEntry } from "./model"; +import { isDefinitiveCanvasConflict } from "./canvas-conflict"; export type ChannelSetup = { canvas: string; agents: string[]; @@ -291,8 +292,23 @@ export function createChannelSetup({ ); const canvas = await operation("canvas", 40100, setup.canvas, [ ["h", p.id], + ["expected-revision", "none"], ]); - await confirm(canvas); + try { + await confirm(canvas); + } catch (error) { + const failed = local + .snapshot() + .find((item) => item.event.id === canvas); + if (isDefinitiveCanvasConflict(failed)) { + // A refused seed must not fence a later manual Canvas save. + await outbox.dismiss(canvas); + throw new Error( + "Canvas changed before the seed was saved; no agents were added by this attempt.", + ); + } + throw error; + } if ((await canvasHead(p.id)) !== canvas) throw new Error( "The seed Canvas is not the selected document; no agents were added by this attempt.", diff --git a/src/features/relay/session.ts b/src/features/relay/session.ts index 5fa46c313..7037076df 100644 --- a/src/features/relay/session.ts +++ b/src/features/relay/session.ts @@ -1252,8 +1252,7 @@ export function createRelaySession( throw new Error("Setup is awaiting exact relay confirmation"); }, refresh: (id, member) => workSessions.refresh(id, { member }, false), - canvasHead: async (id) => - (await channelKit.canvas.read(id, { consistency: "strong" }))?.id, + canvasHead: async (id) => (await channelKit.canvas.read(id))?.id, async preflight(input) { if ( !canonicalDetailsName(input.name) || diff --git a/src/features/relay/socket-requests.ts b/src/features/relay/socket-requests.ts index ee1a9cff5..31bd53a3c 100644 --- a/src/features/relay/socket-requests.ts +++ b/src/features/relay/socket-requests.ts @@ -107,11 +107,12 @@ export function createSocketPublications(wake: () => void) { /^(invalid|blocked|restricted|auth-required|rate-limited):/.test( data[3], ) || - // Buzz workflow CAS/authority and artifact CAS refusals precede + // Buzz workflow CAS/authority and Canvas/artifact CAS refusals precede // domain mutation. ([30620, 46020].includes(job.event.kind) && /^(conflict|forbidden):/.test(data[3])) || - (job.event.kind === 45010 && data[3].startsWith("conflict:")); + ([40100, 45010].includes(job.event.kind) && + data[3].startsWith("conflict:")); job.finish( undefined, new SocketRequestError(