From 0c85846ff884d1019122cff3a7959238e5293fb6 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Fri, 2 Oct 2026 09:58:08 -0700 Subject: [PATCH] refactor(reef): stop reconstructing historical delivery bindings (#163205) --- docs/channels/reef.md | 4 + extensions/reef/src/flow-receipts.test.ts | 338 ++++------------------ extensions/reef/src/flow.ts | 130 +-------- extensions/reef/src/trust-store.ts | 4 +- 4 files changed, 71 insertions(+), 405 deletions(-) diff --git a/docs/channels/reef.md b/docs/channels/reef.md index 7dfa980a1d4d..9241333a1397 100644 --- a/docs/channels/reef.md +++ b/docs/channels/reef.md @@ -211,6 +211,10 @@ openclaw message send --channel reef --target @friend --message "hello from my c A send never fails silently. Local guard or relay errors fail the send immediately. Replies and peer guard rejections come back through the flows below. If the peer's claw confirms nothing for about 10 minutes, the sending agent receives a delivery-delay notice. A follow-up arrives once the message is finally delivered or rejected. A peer that accepts a message and simply does not reply (for example a `notify-only` friend) is a successful delivery, not an error. +When upgrading from early Reef versions, sends that were already in flight may keep an unknown delivery status. Reef leaves their protocol journal intact and handles new sends through current delivery records. Check with the friend before manually resending a pre-upgrade message whose status remains unknown. + +A late rejection for one of those historical sends no longer starts the peer's 15-minute rejection cooldown. A later rejection may therefore allow one automatic rephrased resend that the old cooldown would have suppressed. Peer trust, keys, and the guard checks on new sends are unchanged. + Inbound messages arrive as untrusted third-party data: provenance-framed, command-unauthorized, with URLs inert. Depending on the friend's autonomy tier, OpenClaw notifies you or sends a bounded guarded reply: | Tier | Behavior | diff --git a/extensions/reef/src/flow-receipts.test.ts b/extensions/reef/src/flow-receipts.test.ts index 3403b94577e9..b196325ef9df 100644 --- a/extensions/reef/src/flow-receipts.test.ts +++ b/extensions/reef/src/flow-receipts.test.ts @@ -5,7 +5,6 @@ import { generateIdentity, sha256Hex, signReceipt, - type AuditEntry, } from "../protocol/index.js"; import { MemoryAuditStore, MemoryReplayStore } from "../protocol/memory-stores.test-support.js"; import { ReefMessageFlow } from "./flow.js"; @@ -23,10 +22,6 @@ import { import { reefPeerIdentity } from "./friend-types.js"; import { processReefInboxEntriesInOrder, ReefReceiptNotifier } from "./owner-notice.js"; import type { ReefTransportClient } from "./transport.js"; -import { - REEF_OUTBOUND_DELIVERY_MAX_ENTRIES, - REEF_OUTBOUND_DELIVERY_TTL_MS, -} from "./trust-store.js"; import type { InboxEntry } from "./types.js"; beforeEach(resetFlowStoresForTests); @@ -108,242 +103,70 @@ describe("ReefMessageFlow delivery receipts", () => { expect(entries).not.toHaveBeenCalled(); }); - it("confirms a recent pre-binding accepted receipt only once", async () => { - const alice = generateIdentity(); - const bob = reefKeys(); - const trusted = trust({ alice: peerTrust(alice) }); - const audit = new MemoryAuditStore(new Uint8Array(32).fill(15)); - const id = "01JZ0000000000000000000127"; - const text = "queued before delivery bindings"; - await composeOutbound({ - id, - from: "bob#1", - to: "alice#1", - body: { text }, - senderSigningSecretKey: bob.signing.secretKey, - recipientEncryptionPublicKey: alice.encryption.publicKey, - guard: guard(allow), - audit, - policyVersion: "v1", - }); - const flow = createFlow({ alice, bob, audit, trusted }); - const receipt = signedReceipt(alice, { - auditHead: "a".repeat(64), - id, - bodyHash: sha256Hex(canonicalBytes({ text })), - status: "accepted", - }); - const entry: InboxEntry = receiptEntry(receipt); - - await expect(flow.processEntries([entry])).resolves.toEqual([]); - await expect(flow.processEntries([{ ...entry, seq: 2 }])).resolves.toEqual([]); - - const events = (await audit.entries()).map((item) => item.event.type); - expect(events.filter((type) => type === "confirm_delivery")).toHaveLength(1); - expect(events.filter((type) => type === "invalid_delivery_receipt")).toHaveLength(1); - expect(trusted.deliveries.has(`alice:${id}`)).toBe(false); - }); - - it("does not let abandoned proposals evict sealed legacy deliveries", async () => { - const alice = generateIdentity(); - const bob = reefKeys(); - const audit = new MemoryAuditStore(new Uint8Array(32).fill(18)); - const id = "01JZ0000000000000000000132"; - const bodyHash = "a".repeat(64); - const ts = Math.floor(Date.now() / 1_000); - const entries: AuditEntry[] = [ - { - event: { seq: 1, ts, type: "proposal", payload: { id, to: "alice#1", bodyHash } }, - prevHash: "", - entryHash: "", - }, - { - event: { - seq: 2, - ts, - type: "proposal", - payload: { id: "abandoned-0", to: "alice#1", bodyHash }, - }, - prevHash: "", - entryHash: "", - }, - ...Array.from({ length: REEF_OUTBOUND_DELIVERY_MAX_ENTRIES - 1 }, (_, index) => ({ - event: { - seq: index + 3, - ts, - type: "proposal", - payload: { id: `abandoned-${index + 1}`, to: "alice#1", bodyHash }, - }, - prevHash: "", - entryHash: "", - })), - { - event: { - seq: REEF_OUTBOUND_DELIVERY_MAX_ENTRIES + 2, - ts, - type: "envelope", - payload: { id }, - }, - prevHash: "", - entryHash: "", - }, - ]; - vi.spyOn(audit, "entries").mockResolvedValueOnce(entries); - const flow = createFlow({ alice, bob, audit }); - const receipt = signedReceipt(alice, { - id, - bodyHash, - status: "accepted", - }); - - await expect(flow.processEntries([receiptEntry(receipt)])).resolves.toEqual([]); - expect( - (await audit.entries()).filter((entry) => entry.event.type === "confirm_delivery"), - ).toHaveLength(1); - }); - - it("anchors legacy recovery retention to envelope sealing", async () => { - const alice = generateIdentity(); - const bob = reefKeys(); - const audit = new MemoryAuditStore(new Uint8Array(32).fill(20)); - const id = "01JZ0000000000000000000135"; - const bodyHash = "a".repeat(64); - const sealedAt = Math.floor(Date.now() / 1_000); - const proposedAt = sealedAt - Math.ceil(REEF_OUTBOUND_DELIVERY_TTL_MS / 1_000) - 1; - const entries: AuditEntry[] = [ - { - event: { - seq: 1, - ts: proposedAt, - type: "proposal", - payload: { id, to: "alice#1", bodyHash }, - }, - prevHash: "", - entryHash: "", - }, - { - event: { seq: 2, ts: sealedAt, type: "envelope", payload: { id } }, - prevHash: "", - entryHash: "", - }, - ]; - vi.spyOn(audit, "entries").mockResolvedValueOnce(entries); - const flow = createFlow({ alice, bob, audit }); - const receipt = signedReceipt(alice, { - id, - bodyHash, - status: "accepted", - }); - - await flow.processEntries([receiptEntry(receipt, 1, sealedAt)]); - - expect( - (await audit.entries()).filter((entry) => entry.event.type === "confirm_delivery"), - ).toHaveLength(1); - }); - - it("expires candidates after a cached legacy index ages out", async () => { - const alice = generateIdentity(); - const bob = reefKeys(); - const trusted = trust({ alice: peerTrust(alice) }); - const audit = new MemoryAuditStore(new Uint8Array(32).fill(19)); - const id = "01JZ0000000000000000000133"; - const missId = "01JZ0000000000000000000134"; - const bodyHash = "a".repeat(64); - const now = Date.now(); - const ts = Math.floor(now / 1_000); - const entries: AuditEntry[] = [ - { - event: { seq: 1, ts, type: "proposal", payload: { id, to: "alice#1", bodyHash } }, - prevHash: "", - entryHash: "", - }, - { - event: { seq: 2, ts, type: "envelope", payload: { id } }, - prevHash: "", - entryHash: "", - }, - ]; - const auditEntries = vi.spyOn(audit, "entries").mockResolvedValueOnce(entries); - const nowSpy = vi.spyOn(Date, "now").mockReturnValue(now); - const flow = createFlow({ alice, bob, audit, trusted }); - const miss = signedReceipt(alice, { - id: missId, - bodyHash, - status: "accepted", - }); - const receipt = signedReceipt(alice, { - auditHead: "c".repeat(64), - id, - bodyHash, - status: "accepted", - }); - - try { - await flow.processEntries([receiptEntry(miss)]); - nowSpy.mockReturnValue(now + REEF_OUTBOUND_DELIVERY_TTL_MS + 1_000); - await flow.processEntries([receiptEntry(receipt, 2)]); - } finally { - nowSpy.mockRestore(); - } - - expect(auditEntries).toHaveBeenCalledOnce(); - expect( - (await audit.entries()).filter((entry) => entry.event.type === "confirm_delivery"), - ).toHaveLength(0); - expect(trusted.deliveries.has(`alice:${id}`)).toBe(false); - }); - - it("surfaces a recent pre-binding rejection as durable stop-only guidance", async () => { - const alice = generateIdentity(); - const bob = reefKeys(); - const trusted = trust({ alice: peerTrust(alice) }); - const audit = new MemoryAuditStore(new Uint8Array(32).fill(16)); - const id = "01JZ0000000000000000000128"; - const text = "queued rejection before delivery bindings"; - await composeOutbound({ - id, - from: "bob#1", - to: "alice#1", - body: { text }, - senderSigningSecretKey: bob.signing.secretKey, - recipientEncryptionPublicKey: alice.encryption.publicKey, - guard: guard(allow), - audit, - policyVersion: "v1", - }); - const flow = createFlow({ alice, bob, audit, trusted }); - const receipt = signedReceipt(alice, { - id, - bodyHash: sha256Hex(canonicalBytes({ text })), - status: "rejected", - category: "guard_deny", - }); - const rejections = await flow.processEntries([receiptEntry(receipt)]); - const notify = vi.fn(async () => {}); - const receiptNotifier = createReceiptNotifier(trusted, notify); - - await receiptNotifier.notifyRejections(rejections); - - expect(rejections).toEqual([ - { + it.each(["accepted", "rejected"] as const)( + "quarantines historical %s receipts without reconstructing delivery or cooldown state", + async (status) => { + const alice = generateIdentity(); + const bob = reefKeys(); + const trusted = trust({ alice: peerTrust(alice) }); + const originalPeer = structuredClone(trusted.values.get("alice")); + const audit = new MemoryAuditStore(new Uint8Array(32).fill(15)); + const id = "01JZ0000000000000000000127"; + const text = "queued before delivery bindings"; + await composeOutbound({ id, - peer: "alice", - recipient: reefPeerIdentity(peerTrust(alice)), + from: "bob#1", + to: "alice#1", + body: { text }, + senderSigningSecretKey: bob.signing.secretKey, + recipientEncryptionPublicKey: alice.encryption.publicKey, + guard: guard(allow), + audit, + policyVersion: "v1", + }); + const historicalEntries = structuredClone(await audit.entries()); + const auditEntries = vi.spyOn(audit, "entries"); + const flow = createFlow({ alice, bob, audit, trusted }); + const notify = vi.fn(async () => {}); + const receiptNotifier = createReceiptNotifier(trusted, notify); + const receipt = signedReceipt(alice, { + id, + bodyHash: sha256Hex(canonicalBytes({ text })), + status, + ...(status === "rejected" ? { category: "guard_deny" } : {}), + }); + + const rejections = await flow.processEntries([receiptEntry(receipt)]); + expect(rejections).toEqual([]); + await receiptNotifier.notifyRejections(rejections); + expect(auditEntries).not.toHaveBeenCalled(); + expect(notify).not.toHaveBeenCalled(); + expect(trusted.deliveries.size).toBe(0); + expect(trusted.rejectionNotices.size).toBe(0); + expect(trusted.values.get("alice")).toEqual(originalPeer); + const entries = await audit.entries(); + expect(entries.slice(0, historicalEntries.length)).toEqual(historicalEntries); + expect(entries.filter((entry) => entry.event.type === "confirm_delivery")).toHaveLength(0); + expect( + entries.filter((entry) => entry.event.type === "invalid_delivery_receipt"), + ).toHaveLength(1); + + const currentId = await flow.send("alice", "sent after the upgrade"); + const currentReceipt = signedReceipt(alice, { + id: currentId, + bodyHash: sha256Hex(canonicalBytes({ text: "sent after the upgrade" })), + status: "rejected", category: "guard_deny", - reservedNotice: { lastRejectionAt: expect.any(Number) }, - }, - ]); - expect(notify).toHaveBeenCalledWith( - expect.objectContaining({ - peer: "alice", - messageId: id, - allowResend: false, - text: expect.stringMatching(/Stop automatic retries/), - }), - ); - expect(trusted.deliveries.has(`alice:${id}`)).toBe(false); - }); + }); + await receiptNotifier.notifyRejections( + await flow.processEntries([receiptEntry(currentReceipt, 2)]), + ); + expect(notify).toHaveBeenCalledOnce(); + expect(notify).toHaveBeenCalledWith( + expect.objectContaining({ messageId: currentId, allowResend: true }), + ); + }, + ); it("surfaces one resend notice even when a later batch receipt is invalid", async () => { const alice = generateIdentity(); @@ -429,43 +252,6 @@ describe("ReefMessageFlow delivery receipts", () => { ).toHaveLength(3); }); - it("does not recover a signed rejection from an unsealed outbound proposal", async () => { - const alice = generateIdentity(); - const bob = reefKeys(); - const trusted = trust({ alice: peerTrust(alice) }); - const audit = new MemoryAuditStore(new Uint8Array(32).fill(12)); - const auditEntries = vi.spyOn(audit, "entries"); - const flow = createFlow({ alice, bob, audit, trusted }); - const id = "01JZ0000000000000000000113"; - const receipt = signedReceipt(alice, { - id, - bodyHash: "a".repeat(64), - status: "rejected", - category: "guard_deny", - }); - await audit.appendEvent("proposal", { - id, - from: "bob#1", - to: "alice#1", - bodyHash: receipt.bodyHash, - }); - - await expect(flow.processEntries([receiptEntry(receipt)])).resolves.toEqual([]); - const otherId = "01JZ0000000000000000000131"; - const otherReceipt = signedReceipt(alice, { - auditHead: "c".repeat(64), - id: otherId, - bodyHash: receipt.bodyHash, - status: "rejected", - category: "guard_deny", - }); - await expect(flow.processEntries([receiptEntry(otherReceipt, 2)])).resolves.toEqual([]); - expect(auditEntries).toHaveBeenCalledOnce(); - const events = (await audit.entries()).map((entry) => entry.event.type); - expect(events).toContain("invalid_delivery_receipt"); - expect(events).not.toContain("confirm_delivery"); - }); - it("binds receipts and automatic resends to the send-time recipient identity", async () => { const alice = generateIdentity(); const rotatedAlice = generateIdentity(); diff --git a/extensions/reef/src/flow.ts b/extensions/reef/src/flow.ts index 60fbd83584d3..489cd71fb4b6 100644 --- a/extensions/reef/src/flow.ts +++ b/extensions/reef/src/flow.ts @@ -1,4 +1,3 @@ -import { asOptionalRecord } from "openclaw/plugin-sdk/string-coerce-runtime"; import { bodyHash as hashMessageBody, composeInbound, @@ -10,8 +9,6 @@ import { InvalidDeliveryReceiptError, parseHandleEpoch, PipelineError, - verifyReceipt, - type AuditEntry, type AuditStore, type GuardAdapter, type ReplayStore, @@ -27,66 +24,9 @@ import { import { reefMessageTextHash } from "./rejection-resend.js"; import { ReefDeliveredStore, ReviewApprovalStore } from "./state.js"; import { ReefInboxEntryParkedError, ReefTransportClient } from "./transport.js"; -import { - REEF_OUTBOUND_DELIVERY_MAX_ENTRIES, - REEF_OUTBOUND_DELIVERY_TTL_MS, - type ReefTrustStore, -} from "./trust-store.js"; +import type { ReefTrustStore } from "./trust-store.js"; import type { InboxEntry, ReefDeliveryRejection, ReefIngressMessage, ReefKeys } from "./types.js"; -interface LegacyDeliveryCandidate { - to: string; - bodyHash: string; - expiresAt: number; -} - -function buildLegacyDeliveryIndex( - entries: readonly AuditEntry[], -): Map { - const oldest = Math.floor((Date.now() - REEF_OUTBOUND_DELIVERY_TTL_MS) / 1_000); - const sealed = new Map(); - const confirmed = new Set(); - const candidates = new Map(); - for (let index = entries.length - 1; index >= 0; index -= 1) { - const entry = entries[index]!; - const payload = asOptionalRecord(entry.event.payload); - if (entry.event.type === "confirm_delivery") { - if (entry.event.ts < oldest) { - continue; - } - const receipt = asOptionalRecord(payload?.receipt); - if (typeof receipt?.id === "string") { - confirmed.add(receipt.id); - sealed.delete(receipt.id); - } - } else if (entry.event.type === "envelope" && typeof payload?.id === "string") { - if (entry.event.ts >= oldest && !confirmed.has(payload.id)) { - sealed.set(payload.id, entry.event.ts); - } - } else if (entry.event.type === "proposal") { - const sealedAt = typeof payload?.id === "string" ? sealed.get(payload.id) : undefined; - if ( - typeof payload?.id !== "string" || - typeof payload.to !== "string" || - typeof payload.bodyHash !== "string" || - sealedAt === undefined - ) { - continue; - } - sealed.delete(payload.id); - candidates.set(payload.id, { - to: payload.to, - bodyHash: payload.bodyHash, - expiresAt: sealedAt * 1_000 + REEF_OUTBOUND_DELIVERY_TTL_MS, - }); - if (candidates.size === REEF_OUTBOUND_DELIVERY_MAX_ENTRIES) { - break; - } - } - } - return candidates; -} - /** Reserves a protocol-valid id before recipient-visible Reef delivery starts. */ export const prepareReefMessageId = createMonotonicUlidFactory(); @@ -118,7 +58,6 @@ export function isPermanentReefOutboundRejection(error: unknown): boolean { } export class ReefMessageFlow { - private legacyDeliveryIndex?: Promise>; // Entry ids whose last processing outcome parked (pending review, guard // outage): their re-polls skip the duplicate durable read observation. private readonly parkedReadIds = new Set(); @@ -248,12 +187,9 @@ export class ReefMessageFlow { if (!receipt) { return undefined; } - let delivery = this.options.trust.outboundDelivery(entry.peer, entry.id); + const delivery = this.options.trust.outboundDelivery(entry.peer, entry.id); if (!delivery) { - delivery = await this.recoverLegacyDelivery(entry); - if (!delivery) { - return this.quarantineReceipt(entry); - } + return this.quarantineReceipt(entry); } try { await confirmDelivery(receipt, delivery.recipient.ed25519PublicKey, this.options.audit, { @@ -261,7 +197,6 @@ export class ReefMessageFlow { bodyHash: delivery.bodyHash, ...(delivery.rejection ? { status: "rejected" as const } : {}), }); - await this.forgetLegacyCandidate(entry.id); if (!matchesReefPeerIdentity(this.options.trust.get(entry.peer), delivery.recipient)) { this.options.trust.discardOutboundDelivery(entry.peer, entry.id, delivery); return undefined; @@ -319,65 +254,6 @@ export class ReefMessageFlow { } } - private async recoverLegacyDelivery( - entry: InboxEntry, - ): Promise> { - const receipt = entry.receipt; - const friend = this.options.trust.get(entry.peer); - if (!receipt || receipt.id !== entry.id || !friend || friend.safetyNumberChanged) { - return undefined; - } - if (!verifyReceipt(receipt, friend.ed25519PublicKey)) { - return undefined; - } - const candidates = await this.loadLegacyDeliveryIndex(); - const candidate = candidates.get(entry.id); - if (candidate && candidate.expiresAt <= Date.now()) { - candidates.delete(entry.id); - return undefined; - } - if ( - !candidate || - candidate.to !== formatHandleEpoch(entry.peer, friend.keyEpoch) || - candidate.bodyHash !== receipt.bodyHash - ) { - return undefined; - } - // Upgrade bridge: envelopes sent before delivery bindings shipped can - // still return receipts. Never grant automatic resend from recovered state. - // Remove after that release is older than both relay-retention windows. - this.options.trust.recordOutboundDelivery( - entry.peer, - entry.id, - { - bodyHash: receipt.bodyHash, - recipient: reefPeerIdentity(friend), - }, - { resendDisabled: true }, - ); - candidates.delete(entry.id); - return this.options.trust.outboundDelivery(entry.peer, entry.id); - } - - private loadLegacyDeliveryIndex(): Promise> { - if (!this.legacyDeliveryIndex) { - const pending = this.options.audit.entries().then(buildLegacyDeliveryIndex); - this.legacyDeliveryIndex = pending; - void pending.catch(() => { - if (this.legacyDeliveryIndex === pending) { - this.legacyDeliveryIndex = undefined; - } - }); - } - return this.legacyDeliveryIndex; - } - - private async forgetLegacyCandidate(id: string): Promise { - if (this.legacyDeliveryIndex) { - (await this.legacyDeliveryIndex).delete(id); - } - } - private async quarantineReceipt(entry: InboxEntry): Promise { // A peer-protocol violation must not poison the relay cursor. Keep any // outbound binding intact so a later valid receipt can still complete it. diff --git a/extensions/reef/src/trust-store.ts b/extensions/reef/src/trust-store.ts index 5483cb6b8a25..d7575bc044e0 100644 --- a/extensions/reef/src/trust-store.ts +++ b/extensions/reef/src/trust-store.ts @@ -19,9 +19,9 @@ import type { ReefDeliveryRejection, ReefRejectionNoticeState, RelayFriend } fro export const REEF_TRUST_STORE_MAX_ENTRIES = 4_096; export const REEF_TRUST_STORE_NAMESPACE = "peer-state"; const REEF_OUTBOUND_DELIVERY_STORE_NAMESPACE = "outbound-deliveries"; -export const REEF_OUTBOUND_DELIVERY_MAX_ENTRIES = 32_768; +const REEF_OUTBOUND_DELIVERY_MAX_ENTRIES = 32_768; const REEF_RELAY_RETENTION_MS = 30 * 24 * 60 * 60 * 1_000; -export const REEF_OUTBOUND_DELIVERY_TTL_MS = REEF_RELAY_RETENTION_MS * 2 + 24 * 60 * 60 * 1_000; +const REEF_OUTBOUND_DELIVERY_TTL_MS = REEF_RELAY_RETENTION_MS * 2 + 24 * 60 * 60 * 1_000; const REEF_PAIRING_APPROVAL_PREFIX = "reef-approval-v1:"; const SHA256_HEX_PATTERN = /^[a-f0-9]{64}$/; const MESSAGE_ID_PATTERN = /^[0-7][0-9A-HJKMNP-TV-Z]{25}$/;