From 2a81cc01e153146b743b6e610ae0abcec7df4590 Mon Sep 17 00:00:00 2001 From: grumbach Date: Tue, 29 Sep 2026 16:30:40 +0900 Subject: [PATCH 1/8] feat(pointer): keep the first final state, and look before taking one A pointer state at u64::MAX is now final (ant-protocol): nothing replaces it, not even another final state whose target sorts first. That is what lets an owner hand a pointer's address over for good, by signing one last state that points at the new owner's pointer. The store, the admission gate, fresh offers, repair and hints all compare with the protocol's replaces(), so they keep whichever final state they took first with no code of their own. The merge rule alone leaves one gap: a node holding no final state takes any, so a former owner could still finalize again on a node that joined the group after the transfer, or lost its copy. So before a node takes a final state it does not hold, from a client or a fresh offer, it asks its close group which state each holds. A peer claiming a different final state is asked for the record, and if it verifies -- only the owner could have signed it -- the write is refused as stale, naming the state the group proved. A claim alone refuses nothing, so one peer cannot block a transfer. The look runs after payment is verified, only asks peers that have sent a pointer message, and is bounded at four seconds inside the client's store timeout; silence proves nothing. Two different final states can still exist if the owner races them to different nodes. Each node keeps its first; the client decides by the close group's majority. So the possession check no longer penalises a member holding the other side of such a fork -- it holds what the merge rule told it to -- and repair, between two final states a wide group backs at quorum, adopts the larger side rather than whichever answered first. with_pointers now takes the service and wires both directions, so the node and the e2e harness cannot attach one and forget the other. ADR-0018 records the decision and amends ADR-0016's merge rule. --- Cargo.lock | 4 +- Cargo.toml | 2 +- docs/adr/ADR-0016-pointers-immutable-owner.md | 17 +- ...8-pointer-transfer-by-final-redirection.md | 250 +++++++++++++++++ docs/adr/README.md | 1 + src/node.rs | 7 +- src/pointer/mod.rs | 22 +- src/pointer/service.rs | 258 +++++++++++++++++- src/replication/mod.rs | 24 +- src/replication/pointer.rs | 241 +++++++++++++++- tests/e2e/pointer_replication.rs | 187 ++++++++++++- tests/e2e/testnet.rs | 8 +- tests/pointer_convergence.rs | 175 ++++++++---- 13 files changed, 1079 insertions(+), 117 deletions(-) create mode 100644 docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md diff --git a/Cargo.lock b/Cargo.lock index 27a074bf..c94a4c29 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -883,7 +883,7 @@ dependencies = [ [[package]] name = "ant-protocol" version = "3.0.0" -source = "git+https://github.com/WithAutonomi/ant-protocol?rev=d11d1010d93958fee6e63e4aace8770f39a0678b#d11d1010d93958fee6e63e4aace8770f39a0678b" +source = "git+https://github.com/WithAutonomi/ant-protocol?rev=6f53bb71c62002a312a373281761fe0fdc0fc3aa#6f53bb71c62002a312a373281761fe0fdc0fc3aa" dependencies = [ "blake3", "bytes", @@ -7103,7 +7103,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.48.0", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 03efec06..cc98b10f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -233,7 +233,7 @@ webrtc-direct = [ # the published 3.0.0, so this is the only entry that has to leave the release # baseline. A rev, not a branch, so the pin is immutable. Drop it once the # pointer PR lands and a release includes it. Versions are bumped at release. -ant-protocol = { git = "https://github.com/WithAutonomi/ant-protocol", rev = "d11d1010d93958fee6e63e4aace8770f39a0678b" } +ant-protocol = { git = "https://github.com/WithAutonomi/ant-protocol", rev = "6f53bb71c62002a312a373281761fe0fdc0fc3aa" } [profile.release] lto = true diff --git a/docs/adr/ADR-0016-pointers-immutable-owner.md b/docs/adr/ADR-0016-pointers-immutable-owner.md index ac1a1909..8c5648ff 100644 --- a/docs/adr/ADR-0016-pointers-immutable-owner.md +++ b/docs/adr/ADR-0016-pointers-immutable-owner.md @@ -3,7 +3,7 @@ - **Status:** Proposed - **Date:** 2026-09-18 - **Decision owners:** Anselme (@grumbach) -- **Related:** ADR-0002 (audit), ADR-0004 (commitment-bound pricing), ADR-0008 (per-record pricing), ADR-0009 (audit families), ADR-0011 (capacity-gated discovery), ADR-0014 (file store), ADR-0015 (browser clients) +- **Related:** ADR-0002 (audit), ADR-0004 (commitment-bound pricing), ADR-0008 (per-record pricing), ADR-0009 (audit families), ADR-0011 (capacity-gated discovery), ADR-0014 (file store), ADR-0015 (browser clients), ADR-0018 (a final state is final, and transfer by final redirection — amends the merge rule below) ## Context @@ -85,8 +85,8 @@ rule either would refuse every later update for good, and three such peers would leave no write able to reach its quorum. Under the merge rule each takes the next update, however far ahead of it that is. -An owner who jumps straight to `u64::MAX` limits only themself, and does not -even freeze the pointer: equal counters still resolve by target, below. +An owner who jumps straight to `u64::MAX` limits only themself: the state is +final (ADR-0018), so it is the last one that pointer will ever hold. ### Merge @@ -95,6 +95,10 @@ even freeze the pointer: equal counters still resolve by target, below. 2. smaller target bytes ``` +ADR-0018 puts one rule ahead of these: a state at `u64::MAX` is replaced by +nothing, so two different final states are unordered and a node keeps the +first it took. Below the final counter what follows holds unchanged. + A total order on the states of **one address**. Records of different owners are not comparable and never contend. **Equal state never replaces**: ML-DSA signing is randomized, so one state has unboundedly many valid encodings, and ordering @@ -245,9 +249,10 @@ verified when it was committed and every read verifies it again. self-proving, so one copy settles it; a pointer read has to decide which of several signed states is current, and that answer has to come from more than one peer. -- **Ownership cannot change.** Handover is indirection: point at a new pointer - the recipient owns. The old owner keeps write access forever, so it is a - revocable forwarding state, not a sale. +- **The owner key cannot change.** Handover is indirection: point at a new + pointer the recipient owns. Below the final counter the old owner can take + that back, so it is a revocable forwarding; ADR-0018 makes it stick by + signing it at the final counter. - **Key compromise is permanent.** No rotation, no recovery. - The inlined key costs 1,952 bytes on every read, forever — a deliberate trade for self-contained validation. diff --git a/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md b/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md new file mode 100644 index 00000000..642bd7df --- /dev/null +++ b/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md @@ -0,0 +1,250 @@ +# ADR-0018: Pointer ownership transfer by final redirection + +- **Status:** Proposed +- **Date:** 2026-09-29 +- **Decision owners:** Anselme (@grumbach) +- **Reviewers:** TBD +- **Supersedes:** none. Amends ADR-0016's merge rule at the final counter. +- **Superseded by:** none +- **Related:** ADR-0016 (pointers), ADR-0005 (repair quorum); WithAutonomi/ant-protocol `feat/pointer-ownership-transfer`, WithAutonomi/ant-client `feat/pointer-ownership-transfer` + +## Context + +ADR-0016 fixes a pointer's owner key at creation and offers handover only as +indirection: the owner points the pointer at a pointer the recipient owns. The +former owner keeps its key, so that is a revocable forwarding and not a sale. +The one thing that could make it stick is the counter running out, and under +ADR-0016 it does not: + +```text +1. larger counter +2. smaller target bytes +``` + +At `counter == u64::MAX` no counter is larger, but rule 2 still applies. A +former owner who handed the address over at `u64::MAX` can sign another state +at `u64::MAX` whose target sorts first — about two tries of grinding — and it +displaces the handover on every node, deterministically. ADR-0016 says so +("does not even freeze the pointer") and tells owners to migrate *before* the +final counter, which is exactly the revocable forwarding again. + +People want to hand pointers over — a name, an application's root, a published +handle — while every reader keeps using the same address. The address derives +from the owner key, so it cannot follow a new key. What the address *resolves +to* can. + +## Decision Drivers + +- The address readers use must not change. +- No new record type, field, message or storage format; nothing a node has to + interpret beyond the counter it already compares. +- Once the network holds a transfer, the former owner must have no move left. +- Forks the owner can still make must be detectable by any reader, and must be + impossible to make once a transfer is established. +- Below the final counter nothing changes: ADR-0016's convergence argument + still holds there. + +## Considered Options + +1. **ADR-0016 as is.** Handover by revocable forwarding. Not a transfer. +2. **Certificates and epochs** (ADR-0016's withdrawn alternative): a transfer + certificate chain, an admission index per epoch and a 5-of-7 branch quorum. + Real re-keying, at the cost of a second record type, lineage walks, Sybil + exposure in ownership and a stuck-not-reversed failure mode under churn. +3. **Final state resolved by target, as today.** Grindable in about two tries. +4. **Final state resolved by payment time or order.** A pre-buy defeats it: + pay early, withhold, publish later. Clock trust besides. +5. **A final state is replaced by nothing; the first one a node takes is the + one it keeps.** Clients settle the owner's only remaining fork by the close + group's majority, and nodes look before taking one. + +## Decision + +We choose option 5. + +### Merge + +One rule ahead of ADR-0016's two: + +```text +0. a final state (counter == u64::MAX) is replaced by nothing +1. larger counter +2. smaller target bytes +``` + +Below the final counter the order is ADR-0016's, total and deterministic. At +it, two *different* final states are unordered, so each node keeps whichever it +took first. `replaces` stays a strict partial order — never both ways round, +transitive — and every pair of distinct states is ordered except two final +ones. + +This lives in `ant_protocol::pointer::PointerState::replaces`, so the node's +store, its admission gate, fresh offers, repair and hints all take it without a +line of their own. + +### Transfer + +A transfer is the final state whose target is another pointer: + +```text +counter = u64::MAX +target = (Pointer, recipient) recipient = the new owner's pointer address +``` + +`Pointer::transfer_to` signs it; `transferred_to()` recognises it. A final state +with any other target freezes the pointer — nobody receives it. A pointer target +below the final counter is still ADR-0016's revocable forwarding. + +Readers already follow pointer targets (`pointer_resolve`), so a transferred +address resolves through the recipient's pointer, which only the recipient can +move. The recipient's pointer should serve that one address: a pointer's +address derives from its owner key, so everything handed to the same recipient +pointer resolves to the same place. A recipient uses a fresh key per received +pointer. + +### The fork that remains + +Only the owner can sign a final state, so only the owner can fork one: sign two +and send them to different nodes before either has heard of the other. Each +node keeps its first. Nothing local can settle that, and this design does not +try. What it guarantees instead: + +- **No fork after the fact.** Once a node holds a final state, no arrival — + client PUT, fresh offer or repair — moves it off it. A second final state can + only land on a node that holds none. +- **Nodes look before a final state.** Before taking a final state it does not + hold, from a client or from a fresh offer, a node asks its close group which + state each holds. A peer claiming a *different* final state is asked for the + record, and if it verifies — only the owner could have signed it — the node + refuses, answering `Stale` with the state the group proved. A claim alone + refuses nothing, so one dishonest peer cannot block a transfer. The look runs + after payment is verified (an owner flooding its own finals pays for each + round trip it causes), asks only peers that have sent a pointer message, and + is bounded at four seconds inside the client's ten-second store timeout; + silence proves nothing and the write goes ahead. This is what closes the gap + the merge rule leaves: a node that joined the group after the transfer, or + lost its copy, would otherwise take a second final state on the merge rule + alone. +- **Readers see the majority, and see forks.** A client read that meets a final + state is settled only once one final state is held by a majority of the close + group, and returns that one. If two different final states are seen and + neither has a majority, the read fails as forked rather than guess. One final + state below a majority with no rival is returned once corroborated, as any + state is: it is a transfer still spreading. +- **Recipients check before they rely on it.** `pointer_finality` asks the + whole group and answers `Open`, `Settling` (one final state, short of a + majority), `Final` (one final state, a majority holds it, no rival seen) or + `Forked` (rivals seen; the majority's state if there is one). A recipient + treats the transfer as done only on `Final`; anything else is the former + owner's to explain. + +### Replication + +- **Repair** is the merge rule over quorum-backed states, so a node holding a + final state adopts nothing else, and a node holding none adopts the final + state a quorum holds. Between two final states a group wide enough to back + both at quorum — never a seven-node group at four — the larger side is + adopted, not whichever answered first. +- **Hints** for a different final state are dropped by a node holding one, + since the hint cannot replace it. No refetch loop. +- **Possession.** A member holding a *different final* state is not penalised + when checked for the one this node offered: it holds what the merge rule told + it to, and the fork is the owner's. It is logged at warn. A member holding + nothing is penalised as before. + +### Client + +- `pointer_transfer` refuses before paying if the pointer is already final, if + the recipient is the pointer itself, if the recipient pointer does not exist, + or if the recipient's chain leads back to this address. It then signs, + pays, stores and reports `pointer_finality`. +- `pointer_update` on a final pointer refuses before paying, as it always did + when the counter could not advance. +- `pointer_controller` follows transfers — final pointer-targeted states only — + to the pointer whose owner now decides what the address resolves to. + +## Consequences + +### Positive + +- Pointers can be handed over for good, and the address readers use never + changes. +- No new record, field, message or storage format. The wire is untouched; the + only protocol change is one comparison. +- After the handover lands the former owner has no move at all: no larger + counter exists and no equal one replaces. +- A handover can be chained: the recipient can transfer its own pointer on. +- ADR-0016's "revocable forwarding" and "migrate before the final counter" + caveats are gone: migration *is* the final update. + +### Negative / Trade-offs + +- **A race is a permanent fork.** An owner who signs two final states and races + them leaves each node on its first, forever. With a majority on one side, + reads return that side and flag the fork to anyone who checks finality. With + no majority, reads of that pointer fail as forked, and nothing — not the + owner, not repair — can fix it. Only the owner can do this, and only to its + own pointer, as only the owner can lose its key. +- **A transfer is irreversible.** One sent to the wrong recipient pointer is + gone. The client checks the recipient exists and does not loop back; it + cannot check intent. +- **The look before a final state is not a lock.** It is fail-open on silence + and bounded in time, and it only runs when a node has learned which peers + understand pointers, which takes a sync round after joining. A node asked in + that first round, or one whose group cannot answer in four seconds, takes a + second final state on the merge rule alone. Reads still return the majority + side; the residual risk is a minority fork that `pointer_finality` reports. +- A final PUT costs its node one state query per capable close-group peer, and + a fetch per claimed rival, before it commits. +- **Mixed fleets.** A node on ADR-0016's rule still lets a smaller-target final + state displace the first. Until a close group's majority runs this rule a + former owner can win back the nodes that do not. Reads follow the majority, so + the handover holds wherever most of the group has upgraded; recipients should + wait for the upgrade before relying on a transfer. +- Transferred reads take one extra hop, and a chain of transfers one per hop, + bounded by the client's resolve depth. + +### Neutral / Operational + +- A node refusing a final state because its group proved another logs it at + info with both state ids; a possession check that finds the other side of a + fork logs at warn. Either is the owner equivocating. +- Pricing, payments, audits and storage commitments are unchanged: a final + state is a state, paid, stored, committed and audited like any other. + +## Validation + +- Protocol: a final state is replaced by no counter and no other final state, + whatever its target; below the final counter the order is unchanged; + `replaces` is a strict partial order over a set including final states, total + except between two of them; a fold keeps the first final state it meets; + `transfer_to` signs at the final counter to the recipient, refuses to sign past + a final record, and only a final pointer target counts as a transfer. +- Node, store: a stored transfer is not displaced by a final state whose target + sorts first, nor by any lower counter. +- Node, properties: with one final state among the records every delivery order + converges on it, exhaustively over all permutations; with two, the first + delivered is kept. +- Node, request handler: a final state the group proves already superseded is + refused and named, and nothing is written; one nobody contradicts is taken; + the group is asked only about a final state the node lacks; two final states + raced to two nodes leave each on its first. +- Node, repair: a node holding a final state adopts nothing else; of two final + states with quorum, the larger side is adopted in either answer order. +- Node, live network: a transfer written to one node reaches the group and no + node takes a second final state; a node that missed the transfer refuses a + different one because a peer serves the one it holds, and still takes the + group's own; the possession check does not penalise the other side of a fork + and still penalises a member holding nothing. +- Client: reads over a group holding a transfer, a majority fork, a no-majority + fork and a transfer still spreading; the finality check for each; transfer + preflight refusals; end to end against a local testnet with real settlement, + a transfer is final, readers are redirected to the recipient's pointer, the + recipient moves it, and the former owner's attempt to take it back is refused + by every node. + +## Notes for AI-assisted work + +AI tools may help draft this ADR, but **must not mark it Accepted without human +review**. Accepted ADRs are immutable: create a new superseding ADR rather than +editing an Accepted ADR. diff --git a/docs/adr/README.md b/docs/adr/README.md index 5061984e..69ad2ed7 100644 --- a/docs/adr/README.md +++ b/docs/adr/README.md @@ -36,3 +36,4 @@ See [`TOOLING.md`](./TOOLING.md) for `adrs`, `adr-kit`, and AI harness setup. - [ADR-0013: Settlement version and pre-payment compatibility](./ADR-0013-settlement-version-and-pre-payment-compatibility.md) - [ADR-0015: Direct browser clients over WebRTC Direct](./ADR-0015-direct-browser-clients-over-webrtc-direct.md) - [ADR-0016: Pointers — paid mutable references with an immutable owner](./ADR-0016-pointers-immutable-owner.md) +- [ADR-0018: Pointer ownership transfer by final redirection](./ADR-0018-pointer-transfer-by-final-redirection.md) diff --git a/src/node.rs b/src/node.rs index 56813244..ada1065d 100644 --- a/src/node.rs +++ b/src/node.rs @@ -272,11 +272,10 @@ impl NodeBuilder { }; // ADR-0016: pointers replicate through the same engine. The PUT handler - // hands each newly stored paid state to it on this channel. + // hands each newly stored paid state to it, and asks it before taking + // a final state (ADR-0018). if let Some(service) = protocol.pointer_service() { - let (writes, fresh_writes) = tokio::sync::mpsc::unbounded_channel(); - service.attach_fresh_writes(writes); - engine.with_pointers(service.store().clone(), fresh_writes); + engine.with_pointers(service); } // ADR-0004: wire the engine's commitment state as the quote generator's diff --git a/src/pointer/mod.rs b/src/pointer/mod.rs index c4787f39..e017dfa7 100644 --- a/src/pointer/mod.rs +++ b/src/pointer/mod.rs @@ -3,10 +3,18 @@ //! Implements `docs/adr/ADR-0016-pointers-immutable-owner.md`. //! //! A pointer is a mutable, owner-signed reference stored at an address derived -//! from the owner's public key. Ownership is fixed at creation: there is no -//! transfer, no lineage, no certificates and no key rotation. That choice is -//! what lets the design be this small — the owner key is inlined in the -//! record, so validating a pointer needs nothing but the pointer. +//! from the owner's public key. The owner key is fixed at creation: there is no +//! lineage, no certificates and no key rotation. That choice is what lets the +//! design be this small — the owner key is inlined in the record, so +//! validating a pointer needs nothing but the pointer. +//! +//! What the address resolves to can still be handed over for good +//! (`docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md`): a state at +//! the final counter is replaced by nothing, so an owner who signs one pointing +//! at the new owner's pointer has no move left. The store gets that from the +//! merge rule; the service adds one look at the close group before taking a +//! final state, so a second one cannot land on a node that had not heard of +//! the first. //! //! # What lives here //! @@ -53,8 +61,8 @@ pub mod store; pub use ant_protocol::pointer::{ pointer_address, state_id_for_body, ParsedPointer, Pointer, PointerError, PointerState, - PointerTarget, PointerTargetKind, DATA_TYPE_POINTER, POINTER_BODY_LEN, POINTER_FORMAT_VERSION, - POINTER_WIRE_LEN, TARGET_WIRE_LEN, + PointerTarget, PointerTargetKind, DATA_TYPE_POINTER, FINAL_COUNTER, POINTER_BODY_LEN, + POINTER_FORMAT_VERSION, POINTER_WIRE_LEN, TARGET_WIRE_LEN, }; -pub use service::PointerService; +pub use service::{FinalStateWitness, PointerService}; pub use store::{PointerStore, PutOutcome}; diff --git a/src/pointer/service.rs b/src/pointer/service.rs index bdda6c54..2578b8e0 100644 --- a/src/pointer/service.rs +++ b/src/pointer/service.rs @@ -16,11 +16,19 @@ //! are both 32 bytes from the same range; the domain separator makes a //! collision infeasible, not impossible. A node that holds one kind at an //! address refuses the other rather than silently choosing. +//! 4. **A look before a final state.** A state at the final counter is +//! replaced by nothing, so a node that takes one can never be corrected. Before +//! taking one it does not already hold, the node asks its close group +//! whether a *different* final state is already held there, and refuses if a +//! peer proves one with the signed record (ADR-0018). That is what keeps a +//! former owner from handing an address over twice, to nodes that had not +//! heard of the first handover yet. //! //! # Order of work //! //! ```text -//! parse → compare with what is held → verify signature → check payment → commit +//! parse → compare with what is held → verify signature → check payment +//! → (final state only) ask the close group → commit //! ``` //! //! The comparison precedes the signature check, so re-submitting a state the @@ -35,6 +43,7 @@ use ant_protocol::chunk::{ XorName, }; use bytes::Bytes; +use futures::future::BoxFuture; use parking_lot::RwLock; use saorsa_core::P2PNode; use tokio::sync::mpsc; @@ -46,7 +55,23 @@ use crate::pointer::store::{Inspected, PointerStore, PutOutcome}; use crate::replication::admission; use crate::replication::pointer::PointerFreshWrite; use crate::storage::{ChunkStore, SELF_CLOSENESS_GATE_WIDTH}; -use ant_protocol::pointer::POINTER_WIRE_LEN; +use ant_protocol::pointer::{Pointer, PointerState, POINTER_WIRE_LEN}; + +/// Where a node looks, before it takes a final state, for a different final +/// state its close group already holds (ADR-0018). +/// +/// A trait rather than the replication engine itself, because the engine is +/// built after this service and needs a running P2P node, and because what the +/// service decides from the answer is worth testing without one. +pub trait FinalStateWitness: Send + Sync { + /// A final state at `state.address`, other than `state`, that a peer in + /// the close group holds — as the signed record, which only the owner can + /// have made — or `None` if none was found. + /// + /// `None` also when the group could not be asked in time. Silence is not + /// evidence, so it never refuses a write; only a record does. + fn conflicting_final<'a>(&'a self, state: &'a PointerState) -> BoxFuture<'a, Option>; +} /// Handles pointer requests against a [`PointerStore`]. #[derive(Clone)] @@ -70,6 +95,10 @@ pub struct PointerService { /// handler before it has a running P2P node. `None` in unit tests that /// never attach one, exactly as the chunk path does. p2p_node: Arc>>>, + /// Asked before a final state is taken (ADR-0018). Attached with the + /// replication engine; `None` where nothing replicates, and then a final + /// state is taken on the merge rule alone. + final_witness: Arc>>>, } impl std::fmt::Debug for PointerService { @@ -92,6 +121,7 @@ impl PointerService { payments: None, fresh_writes: Arc::new(RwLock::new(None)), p2p_node: Arc::new(RwLock::new(None)), + final_witness: Arc::new(RwLock::new(None)), } } @@ -125,6 +155,11 @@ impl PointerService { *self.fresh_writes.write() = Some(writes); } + /// Ask `witness` before taking a final state this node does not hold. + pub fn attach_final_state_witness(&self, witness: Arc) { + *self.final_witness.write() = Some(witness); + } + /// The store this service fronts. #[must_use] pub const fn store(&self) -> &PointerStore { @@ -162,10 +197,11 @@ impl PointerService { } }; - let address = parsed.state().address; - let state_id = parsed.state().state_id; + let state = *parsed.state(); + let address = state.address; + let state_id = state.state_id; - if let Some(refusal) = self.admit(parsed.state()).await { + if let Some(refusal) = self.admit(&state).await { return refusal; } @@ -191,6 +227,17 @@ impl PointerService { } } + // Last, because it costs the close group a round trip, and only for a + // paid final state: if a peer proves a different final state is + // already held, this one lost the race to it everywhere that matters. + // Named as the state this node knows of, as for any stale arrival. + if let Some(conflict) = self.conflicting_final(&state).await { + return PointerPutResponse::Stale { + address, + state_id: conflict.state_id(), + }; + } + // Charge the bytes this write will take, and hold the charge until it // lands. Checking capacity and then writing is the race the file store // exists to close: concurrent writers all pass one cached measurement @@ -251,10 +298,7 @@ impl PointerService { /// `Some(response)` means refuse. Ordered cheapest first, and all of it /// ahead of the signature check, so a forged record for an address this /// node does not serve buys no cryptography. - async fn admit( - &self, - state: &ant_protocol::pointer::PointerState, - ) -> Option { + async fn admit(&self, state: &PointerState) -> Option { let address = state.address; // Any state the merge rule prefers to what this node knows, whatever @@ -319,6 +363,34 @@ impl PointerService { None } + /// The different final state the close group proves it already holds, if + /// `state` is a final state this node would be taking for the first time. + /// + /// Not asked for anything else. A lower counter can be replaced, so taking + /// one blind costs nothing a later state cannot fix; and a node that + /// already holds a final state has had its answer from the merge rule — + /// the same state is unchanged and any other is stale — before this runs. + async fn conflicting_final(&self, state: &PointerState) -> Option { + if !state.is_terminal() + || self + .store + .state(&state.address) + .is_some_and(|held| held.is_terminal()) + { + return None; + } + let witness = self.final_witness.read().as_ref().map(Arc::clone)?; + let conflict = witness.conflicting_final(state).await?; + info!( + "Refusing final pointer state {} at {}: the close group already holds final \ + state {}, so the owner has finalized it before", + hex::encode(state.state_id), + hex::encode(state.address), + hex::encode(conflict.state_id()) + ); + Some(conflict) + } + /// Handle a pointer GET. pub async fn handle_get(&self, request: PointerGetRequest) -> PointerGetResponse { match self.store.get(&request.address).await { @@ -781,6 +853,174 @@ mod tests { assert_eq!(held.counter(), u64::MAX); } + /// A close group whose answer to the finality question is fixed, and + /// which counts how often it was asked. + struct StubWitness { + conflict: Option, + asked: std::sync::atomic::AtomicUsize, + } + + impl StubWitness { + fn new(conflict: Option) -> Arc { + Arc::new(Self { + conflict, + asked: std::sync::atomic::AtomicUsize::new(0), + }) + } + + fn asked(&self) -> usize { + self.asked.load(std::sync::atomic::Ordering::SeqCst) + } + } + + impl FinalStateWitness for StubWitness { + fn conflicting_final<'a>( + &'a self, + _state: &'a PointerState, + ) -> BoxFuture<'a, Option> { + self.asked.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + let conflict = self.conflict.clone(); + Box::pin(async move { conflict }) + } + } + + fn final_state(seed: u8, target_byte: u8) -> Pointer { + let (pk, sk) = keypair(seed); + let target = PointerTarget::new(PointerTargetKind::Pointer, [target_byte; 32]); + Pointer::sign(&sk, &pk, ant_protocol::pointer::FINAL_COUNTER, target).expect("sign") + } + + #[tokio::test] + async fn a_final_state_the_group_proves_is_already_superseded_is_refused() { + // The node has the pointer's ordinary state and has never seen a final + // one, so the merge rule alone would take this. The group proves the + // owner already finalized it elsewhere, so it is refused, named as the + // state that got there first, and nothing is written. + let (service, _dir) = service().await; + let current = signed(1, 3, 1); + service.handle_put(put(¤t)).await; + + let established = final_state(1, 0xAA); + let late = final_state(1, 0x01); + let witness = StubWitness::new(Some(established.clone())); + service.attach_final_state_witness(witness.clone()); + + match service.handle_put(put(&late)).await { + PointerPutResponse::Stale { address, state_id } => { + assert_eq!(address, late.address()); + assert_eq!( + state_id, + established.state_id(), + "the first final state is named" + ); + } + other => panic!("a second final state must be refused, got {other:?}"), + } + assert_eq!(witness.asked(), 1); + assert_eq!( + service.store().state_id(&late.address()), + Some(current.state_id()), + "nothing was written" + ); + } + + #[tokio::test] + async fn a_final_state_nobody_contradicts_is_taken() { + // Silence is not evidence: a group that proves nothing lets the final + // state through on the merge rule. + let (service, _dir) = service().await; + let witness = StubWitness::new(None); + service.attach_final_state_witness(witness.clone()); + + let transfer = final_state(2, 0x42); + assert!(matches!( + service.handle_put(put(&transfer)).await, + PointerPutResponse::Success { .. } + )); + assert_eq!(witness.asked(), 1); + assert_eq!( + service.store().state_id(&transfer.address()), + Some(transfer.state_id()) + ); + } + + #[tokio::test] + async fn the_group_is_asked_only_about_a_final_state_the_node_lacks() { + // A lower counter can always be replaced later, so it is taken without + // a round trip. A node already holding a final state has its answer + // from the merge rule before any round trip: the same state is + // unchanged, and any other is stale. + let (service, _dir) = service().await; + let witness = StubWitness::new(None); + service.attach_final_state_witness(witness.clone()); + + for counter in [0u64, 1, u64::MAX - 1] { + assert!(matches!( + service.handle_put(put(&signed(3, counter, 1))).await, + PointerPutResponse::Success { .. } + )); + } + assert_eq!(witness.asked(), 0, "no final state yet, no question"); + + let first = final_state(3, 0x10); + assert!(matches!( + service.handle_put(put(&first)).await, + PointerPutResponse::Success { .. } + )); + assert_eq!(witness.asked(), 1); + + let second = final_state(3, 0x01); + match service.handle_put(put(&second)).await { + PointerPutResponse::Stale { state_id, .. } => { + assert_eq!(state_id, first.state_id()); + } + other => panic!("a second final state must be stale, got {other:?}"), + } + assert!(matches!( + service.handle_put(put(&final_state(3, 0x10))).await, + PointerPutResponse::Unchanged { .. } + )); + assert_eq!(witness.asked(), 1, "a held final state answers for itself"); + } + + #[tokio::test] + async fn two_final_states_raced_to_two_nodes_leave_each_on_its_first() { + // The one fork ADR-0018 allows: the owner signs two final states and + // sends one to each node before either has heard of the other. Each + // keeps what it took first and refuses the other, whichever sorts + // first. Nothing on a node can settle it; a reader decides by how many + // of the group hold each side. + let (first_node, _a) = service().await; + let (second_node, _b) = service().await; + let one = final_state(4, 0x09); + let other = final_state(4, 0x01); + + assert!(matches!( + first_node.handle_put(put(&one)).await, + PointerPutResponse::Success { .. } + )); + assert!(matches!( + second_node.handle_put(put(&other)).await, + PointerPutResponse::Success { .. } + )); + assert!(matches!( + first_node.handle_put(put(&other)).await, + PointerPutResponse::Stale { .. } + )); + assert!(matches!( + second_node.handle_put(put(&one)).await, + PointerPutResponse::Stale { .. } + )); + assert_eq!( + first_node.store().state_id(&one.address()), + Some(one.state_id()) + ); + assert_eq!( + second_node.store().state_id(&other.address()), + Some(other.state_id()) + ); + } + #[test] fn a_pointer_is_refused_where_a_chunk_already_sits() { // The two kinds share one 32-byte address space. A collision is diff --git a/src/replication/mod.rs b/src/replication/mod.rs index 73d91f45..6e877d17 100644 --- a/src/replication/mod.rs +++ b/src/replication/mod.rs @@ -68,6 +68,7 @@ use crate::payment::{ MIN_PAYMENT_PROOF_SIZE_BYTES, }; use crate::pointer::store::PointerStore; +use crate::pointer::{FinalStateWitness, PointerService}; use crate::replication::audit::AuditTickResult; use crate::replication::audit_coordinator::AuditChallengeCoordinator; use crate::replication::audit_metrics::{ @@ -2214,17 +2215,17 @@ impl ReplicationEngine { }) } - /// Replicate pointers too (ADR-0016): the records in `store`, and the - /// fresh writes the pointer PUT handler sends on `fresh_writes`. + /// Replicate the pointers `service` stores (ADR-0016), and wire the + /// service to this engine both ways: it hands each newly stored paid + /// state here to be offered on, and asks here before it takes a final + /// state (ADR-0018). /// /// Call before [`Self::start`]. - pub fn with_pointers( - &mut self, - store: PointerStore, - fresh_writes: mpsc::UnboundedReceiver, - ) { - self.pointers = Some(Arc::new(pointer::PointerReplication::new( - store, + pub fn with_pointers(&mut self, service: &PointerService) { + let (writes, fresh_writes) = mpsc::unbounded_channel(); + service.attach_fresh_writes(writes); + let replication = Arc::new(pointer::PointerReplication::new( + service.store().clone(), Arc::clone(&self.storage), Arc::clone(&self.p2p_node), Arc::clone(&self.payment_verifier), @@ -2233,7 +2234,10 @@ impl ReplicationEngine { Arc::clone(&self.send_semaphore), self.shutdown.clone(), self.detached_task_tracker.clone(), - ))); + )); + let witness: Arc = replication.clone(); + service.attach_final_state_witness(witness); + self.pointers = Some(replication); self.pointer_fresh_rx = Some(fresh_writes); } diff --git a/src/replication/pointer.rs b/src/replication/pointer.rs index df62a0f4..1fb49aaa 100644 --- a/src/replication/pointer.rs +++ b/src/replication/pointer.rs @@ -24,7 +24,13 @@ //! they hold that state or a newer one by returning a valid record. //! - **Possession.** Some minutes after offering a fresh state, the node asks //! each close-group member for the record. A member that is still responsible -//! and cannot produce that state or a newer one is penalised. +//! and cannot produce that state or a newer one is penalised — unless it holds +//! a *different final* state, which is a fork its owner made and not a failure +//! to store. +//! - **Finality.** Before taking a final state it does not hold, from a client +//! or a fresh offer, a node asks the close group whether a different final +//! state is already held, and refuses if a peer proves one with the signed +//! record (ADR-0018). See [`PointerReplication::conflicting_final`]. //! //! Requests go only to peers that have sent a pointer message themselves (see //! [`PointerReplication::is_capable`]). A peer built before pointers cannot @@ -38,6 +44,7 @@ use std::sync::Arc; use std::time::{Duration, Instant}; use ant_protocol::pointer::{Pointer, PointerState, POINTER_WIRE_LEN}; +use futures::future::BoxFuture; use futures::stream::{self, StreamExt}; use parking_lot::Mutex; use rand::Rng; @@ -55,6 +62,7 @@ use crate::payment::{ MIN_PAYMENT_PROOF_SIZE_BYTES, }; use crate::pointer::store::{Inspected, PointerStore}; +use crate::pointer::FinalStateWitness; use crate::replication::admission; use crate::replication::commitment_state::ResponderCommitmentState; use crate::replication::config::{ @@ -102,6 +110,15 @@ const UNDECIDED_BACKOFF: Duration = Duration::from_secs(30); /// Most prune candidates examined per pass. const MAX_PRUNE_CANDIDATES_PER_PASS: usize = 256; +/// How long a node spends asking its close group for a conflicting final +/// state before taking one (ADR-0018). +/// +/// It runs inside a client's PUT, after payment has been verified, so it has +/// to leave that PUT well inside the client's ten-second store timeout. A +/// group that has not answered by then is taken to hold nothing that +/// conflicts: silence is never a vote, here as anywhere else. +pub(crate) const FINAL_STATE_CHECK_BUDGET: Duration = Duration::from_secs(4); + /// One peer's answer about one address: the state it holds there, if any. type StateAnswer = ((PeerId, XorName), Option); @@ -171,13 +188,22 @@ pub(crate) fn evaluate( } let beats_held = |state: &PointerState| held.is_none_or(|held| state.replaces(held)); + // The merge winner, and between two final states — which the merge rule + // leaves unordered — the one more of the group holds. Only a group wider + // than twice the quorum can back two at once, and then arrival order must + // not be what decides. + let prefer = |candidate: &(PointerState, Vec), + current: &(PointerState, Vec)| { + candidate.0.replaces(¤t.0) + || (!current.0.replaces(&candidate.0) && candidate.1.len() > current.1.len()) + }; let best = by_state .iter() .filter(|(state, holders)| holders.len() >= quorum_needed && beats_held(state)) .fold( None::<&(PointerState, Vec)>, |best, candidate| match best { - Some(current) if !candidate.0.replaces(¤t.0) => Some(current), + Some(current) if !prefer(candidate, current) => Some(current), _ => Some(candidate), }, ); @@ -953,6 +979,26 @@ impl PointerReplication { ); return; } + // A final state held nowhere here yet is looked for in the group + // first, exactly as a client PUT of one is: the peer offering it may + // simply be the side of a race this node has not heard the other + // side of. + let holds_final = self + .store + .state(&state.address) + .is_some_and(|held| held.is_terminal()); + if state.is_terminal() && !holds_final { + if let Some(conflict) = self.conflicting_final(&state).await { + info!( + "Refusing fresh final pointer state {} at {} from {source}: the close \ + group already holds final state {}", + hex::encode(state.state_id), + hex::encode(state.address), + hex::encode(conflict.state_id()) + ); + return; + } + } match self .store_verified(record, Some(self.config.paid_list_close_group_size)) .await @@ -969,6 +1015,104 @@ impl PointerReplication { } } + // ----------------------------------------------------------------------- + // Finality + // ----------------------------------------------------------------------- + + /// A final state at `state.address`, other than `state`, that a peer in + /// the close group holds and proves by serving the signed record. + /// + /// Asked before this node takes a final state it does not hold. A final + /// state is replaced by nothing, so a node that took a second one would + /// hold it for good; asking first is what keeps a former owner from + /// finalizing an address again on nodes that had not yet heard it was + /// final — one that joined the group since, or lost its copy. + /// + /// A peer's word is not enough to refuse: a state summary is a claim + /// anyone can make, and one dishonest peer could otherwise block every + /// handover. A signed record is not a claim. Only the owner can sign a + /// final state, so one that verifies is the owner's own proof that it + /// finalized the pointer before. + /// + /// Bounded by [`FINAL_STATE_CHECK_BUDGET`], and only peers that have sent + /// a pointer message are asked. A group that cannot be asked in time finds + /// nothing, and the write goes ahead on the merge rule alone: a race is a + /// fork the client detects, not one this can prevent. + pub async fn conflicting_final(&self, state: &PointerState) -> Option { + if !state.is_terminal() { + return None; + } + tokio::time::timeout(FINAL_STATE_CHECK_BUDGET, self.find_conflicting_final(state)) + .await + .unwrap_or_else(|_| { + debug!( + "Close group of pointer {} did not answer the finality check in time", + hex::encode(state.address) + ); + None + }) + } + + /// The body of [`Self::conflicting_final`], without its time bound. + /// + /// Answers are taken as they arrive, and a claim is checked the moment it + /// is made, so one slow peer cannot hide a quick one's proof behind the + /// budget. + async fn find_conflicting_final(&self, state: &PointerState) -> Option { + let self_id = *self.p2p.peer_id(); + let address = state.address; + let peers: Vec = self + .p2p + .dht_manager() + .find_closest_nodes_local(&address, self.config.close_group_size) + .await + .into_iter() + .map(|node| node.peer_id) + .filter(|peer| *peer != self_id && self.is_capable(peer)) + .collect(); + + let mut claims = stream::iter(peers) + .map(|peer| async move { + let held = self.ask_state(&peer, &address).await?; + (held.is_terminal() && held.state_id != state.state_id).then_some(peer) + }) + .buffer_unordered(VERIFICATION_CONCURRENCY); + while let Some(claim) = claims.next().await { + let Some(peer) = claim else { continue }; + let Some(record) = self.fetch_record(&peer, &address).await else { + continue; + }; + if record.is_terminal() && record.state_id() != state.state_id { + return Some(record); + } + } + None + } + + /// Ask one peer which state it holds at `address`. + async fn ask_state(&self, peer: &PeerId, address: &XorName) -> Option { + let body = self + .request( + peer, + ReplicationMessageBody::PointerStateRequest(PointerStateRequest { + addresses: vec![*address], + }), + FINAL_STATE_CHECK_BUDGET, + ) + .await?; + let ReplicationMessageBody::PointerStateResponse(response) = body else { + return None; + }; + // An answer about a different address is no answer. + response + .states + .into_iter() + .next() + .flatten() + .filter(|summary| summary.address == *address) + .map(PointerState::from) + } + // ----------------------------------------------------------------------- // Possession // ----------------------------------------------------------------------- @@ -1012,19 +1156,32 @@ impl PointerReplication { if *peer == self_id || !group.contains(peer) || !self.is_capable(peer) { continue; } - let holds = self + let held = self .fetch_record(peer, &fresh.address) .await - .is_some_and(|record| { - let state = record.state(); - state.state_id == fresh.state_id || state.replaces(&fresh) - }); - if !holds { - warn!( - "Peer {peer} does not hold pointer {} it was offered", - hex::encode(fresh.address) - ); - self.penalise(peer).await; + .map(|record| record.state()); + match held { + Some(state) if state.state_id == fresh.state_id || state.replaces(&fresh) => {} + // Two different final states: the owner signed both and each + // node kept the one it took first, as this one did. The peer + // is holding what the merge rule told it to hold, so it is not + // penalised for the owner's fork. Said loudly, because a read + // of this pointer now depends on which side most of the group + // is on. + Some(state) if state.is_terminal() && fresh.is_terminal() => warn!( + "Pointer {} is forked: peer {peer} holds final state {}, this node \ + offered final state {}", + hex::encode(fresh.address), + hex::encode(state.state_id), + hex::encode(fresh.state_id) + ), + _ => { + warn!( + "Peer {peer} does not hold pointer {} it was offered", + hex::encode(fresh.address) + ); + self.penalise(peer).await; + } } } } @@ -1160,6 +1317,12 @@ impl PointerReplication { } } +impl FinalStateWitness for PointerReplication { + fn conflicting_final<'a>(&'a self, state: &'a PointerState) -> BoxFuture<'a, Option> { + Box::pin(Self::conflicting_final(self, state)) + } +} + #[cfg(test)] #[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] mod tests { @@ -1297,6 +1460,58 @@ mod tests { } } + #[test] + fn a_node_holding_a_final_state_adopts_nothing_else() { + // Not even another final state the whole group holds: the merge rule + // replaces a final state with nothing, and repair is the merge rule + // applied to what the group says. + let held = state(u64::MAX, 9, 90); + let other = state(u64::MAX, 1, 91); + let got = evaluate( + Some(&held), + 7, + &answers(&[ + (1, Some(other)), + (2, Some(other)), + (3, Some(other)), + (4, Some(other)), + (5, Some(other)), + ]), + 4, + ); + assert_eq!(got, Verdict::Refused); + } + + #[test] + fn of_two_final_states_with_quorum_the_one_more_hold_is_adopted() { + // Two final states are unordered, so only a group wide enough to back + // both can get here. Arrival order must not pick; the larger side does. + let fewer = state(u64::MAX, 1, 70); + let more = state(u64::MAX, 9, 71); + for order in [[fewer, more], [more, fewer]] { + let mut list = Vec::new(); + let mut next = 1u8; + for candidate in order { + let holders = if candidate.state_id == more.state_id { + 5 + } else { + 4 + }; + for _ in 0..holders { + list.push((next, Some(candidate))); + next += 1; + } + } + match evaluate(None, 20, &answers(&list), 4) { + Verdict::Adopt { state, holders } => { + assert_eq!(state.state_id, more.state_id); + assert_eq!(holders.len(), 5); + } + other => panic!("expected the larger side, got {other:?}"), + } + } + } + #[test] fn a_group_nobody_can_be_asked_in_is_undecided() { assert_eq!(evaluate(None, 0, &[], 0), Verdict::Undecided); diff --git a/tests/e2e/pointer_replication.rs b/tests/e2e/pointer_replication.rs index c031d719..eaa9d449 100644 --- a/tests/e2e/pointer_replication.rs +++ b/tests/e2e/pointer_replication.rs @@ -19,7 +19,9 @@ use ant_node::replication::commitment::pointer_leaf_hash; use ant_node::replication::commitment_state::{BuiltCommitment, ResponderCommitmentState}; use ant_node::replication::pointer::{PointerFreshWrite, PointerReplication}; use ant_node::ReplicationConfig; -use ant_protocol::pointer::{Pointer, PointerState, PointerTarget, PointerTargetKind}; +use ant_protocol::pointer::{ + Pointer, PointerState, PointerTarget, PointerTargetKind, FINAL_COUNTER, +}; use bytes::Bytes; use saorsa_core::identity::PeerId; use saorsa_pqc::api::sig::{ml_dsa_65, MlDsaPublicKey, MlDsaSecretKey}; @@ -51,6 +53,18 @@ fn signed(pk: &MlDsaPublicKey, sk: &MlDsaSecretKey, counter: u64, target: u8) -> .expect("sign") } +/// The final state that hands `pk`'s pointer over to the pointer at +/// `recipient` (ADR-0018). +fn transfer(pk: &MlDsaPublicKey, sk: &MlDsaSecretKey, recipient: u8) -> Pointer { + Pointer::sign( + sk, + pk, + FINAL_COUNTER, + PointerTarget::new(PointerTargetKind::Pointer, [recipient; 32]), + ) + .expect("sign") +} + fn store(node: &TestNode) -> PointerStore { node.ant_protocol .as_ref() @@ -452,6 +466,177 @@ async fn the_possession_check_penalises_only_a_member_that_dropped_the_record() harness.teardown().await.expect("teardown"); } +/// A transfer written to one node reaches the whole group, and after that no +/// node takes a second final state — not even one whose target sorts first, +/// which the previous merge rule let displace the first everywhere (ADR-0018). +#[tokio::test] +#[serial] +async fn a_transfer_reaches_the_group_and_no_node_takes_a_second_one() { + let harness = TestHarness::setup_minimal().await.expect("setup"); + harness.warmup_dht().await.expect("warmup"); + + let (pk, sk) = owner(); + let created = signed(&pk, &sk, 0, 1); + let handed_over = transfer(&pk, &sk, 0x77); + let take_back = transfer(&pk, &sk, 0x01); + assert!(take_back.target().to_bytes() < handed_over.target().to_bytes()); + for record in [&created, &handed_over, &take_back] { + mark_paid(&harness, record); + } + + put(harness.test_node(1).expect("node"), &created).await; + for i in 0..harness.node_count() { + assert!( + wait_for(&harness, i, &created).await, + "node {i} lacks the create" + ); + } + match put(harness.test_node(2).expect("node"), &handed_over).await { + PointerPutResponse::Success { state_id, .. } => { + assert_eq!(state_id, handed_over.state_id()); + } + other => panic!("the transfer was refused: {other:?}"), + } + for i in 0..harness.node_count() { + assert!( + wait_for(&harness, i, &handed_over).await, + "node {i} never received the transfer" + ); + } + + // The former owner, paid up, tries to take it back on every node. + for i in 0..harness.node_count() { + match put(harness.test_node(i).expect("node"), &take_back).await { + PointerPutResponse::Stale { state_id, .. } => assert_eq!( + state_id, + handed_over.state_id(), + "node {i} must name the transfer it holds" + ), + other => panic!("node {i} took a second final state: {other:?}"), + } + assert!(holds(harness.test_node(i).expect("node"), &handed_over)); + } + + harness.teardown().await.expect("teardown"); +} + +/// A node that never heard of a transfer — one that joined after it, or lost +/// its copy — would take any final state on the merge rule alone. Before it +/// does, it asks its group, and a peer that serves the transfer it holds is +/// proof enough to refuse the second one. +#[tokio::test] +#[serial] +async fn a_node_that_missed_the_transfer_refuses_another_its_group_proves() { + let harness = TestHarness::setup_minimal().await.expect("setup"); + harness.warmup_dht().await.expect("warmup"); + + let (pk, sk) = owner(); + let created = signed(&pk, &sk, 0, 1); + let handed_over = transfer(&pk, &sk, 0x77); + let second = transfer(&pk, &sk, 0x01); + mark_paid(&harness, &second); + mark_paid(&harness, &handed_over); + + // Everyone holds the create, and learns everyone else understands + // pointers while there is nothing newer to hint. + let everyone: Vec = (0..harness.node_count()).collect(); + for i in &everyone { + store(harness.test_node(*i).expect("node")) + .put_bytes(&created.to_bytes()) + .await + .expect("put"); + } + exchange_hints(&harness, &everyone, &everyone).await; + + // The transfer lands everywhere but one node, written straight into the + // stores so nothing replicates or hints it there. + let unaware = 4; + for i in others(&harness, &[unaware]) { + store(harness.test_node(i).expect("node")) + .put_bytes(&handed_over.to_bytes()) + .await + .expect("put"); + } + let node = harness.test_node(unaware).expect("node"); + assert!( + holds(node, &created), + "the unaware node has only the create" + ); + + match put(node, &second).await { + PointerPutResponse::Stale { state_id, .. } => assert_eq!( + state_id, + handed_over.state_id(), + "the refusal names the transfer the group proved" + ), + other => panic!("the unaware node took a second final state: {other:?}"), + } + assert!(holds(node, &created), "nothing was written"); + + // The transfer the group holds is not refused: it is no conflict. + match put(node, &handed_over).await { + PointerPutResponse::Success { state_id, .. } => { + assert_eq!(state_id, handed_over.state_id()); + } + other => panic!("the group's own transfer was refused: {other:?}"), + } + + harness.teardown().await.expect("teardown"); +} + +/// Two different final states are a fork only the owner can make, by racing +/// them. A member holding the other side took what reached it first, as the +/// merge rule says; the possession check does not penalise it for the owner's +/// fork, while it still penalises a member that holds nothing. +#[tokio::test] +#[serial] +async fn the_possession_check_does_not_penalise_the_other_side_of_a_fork() { + let harness = TestHarness::setup_minimal().await.expect("setup"); + harness.warmup_dht().await.expect("warmup"); + + let (pk, sk) = owner(); + let one_side = transfer(&pk, &sk, 0x77); + let other_side = transfer(&pk, &sk, 0x01); + let checker = 3; + let forked = 1; + let dropper = 2; + for i in others(&harness, &[forked, dropper]) { + store(harness.test_node(i).expect("node")) + .put_bytes(&one_side.to_bytes()) + .await + .expect("put"); + } + store(harness.test_node(forked).expect("node")) + .put_bytes(&other_side.to_bytes()) + .await + .expect("put"); + + let everyone: Vec = (0..harness.node_count()).collect(); + exchange_hints(&harness, &everyone, &[checker]).await; + + let checker_node = harness.test_node(checker).expect("node"); + let checker_p2p = checker_node.p2p_node.as_ref().expect("p2p"); + let forked_peer = peer(harness.test_node(forked).expect("node")); + let dropper_peer = peer(harness.test_node(dropper).expect("node")); + let forked_before = checker_p2p.peer_trust(&forked_peer); + let dropper_before = checker_p2p.peer_trust(&dropper_peer); + + replication(checker_node) + .check_possession(one_side.state(), &[forked_peer, dropper_peer]) + .await; + + assert!( + checker_p2p.peer_trust(&forked_peer) >= forked_before, + "the member holding the other final state was penalised" + ); + assert!( + checker_p2p.peer_trust(&dropper_peer) < dropper_before, + "the member holding nothing was not penalised" + ); + + harness.teardown().await.expect("teardown"); +} + /// A network whose close group is two nodes, so that in five nodes some node /// is always outside the retention width (two plus the margin of two) for any /// address, with no pruning hysteresis. diff --git a/tests/e2e/testnet.rs b/tests/e2e/testnet.rs index 6e878709..64d7c5e3 100644 --- a/tests/e2e/testnet.rs +++ b/tests/e2e/testnet.rs @@ -1371,11 +1371,11 @@ impl TestNetwork { .await { Ok(mut engine) => { - // Pointers replicate through the same engine (ADR-0016). + // Pointers replicate through the same engine (ADR-0016), + // which the service also asks before a final state + // (ADR-0018). if let Some(service) = protocol.pointer_service() { - let (writes, fresh_writes) = tokio::sync::mpsc::unbounded_channel(); - service.attach_fresh_writes(writes); - engine.with_pointers(service.store().clone(), fresh_writes); + engine.with_pointers(service); } let dht_events = p2p.dht_manager().subscribe_events(); engine.start(dht_events); diff --git a/tests/pointer_convergence.rs b/tests/pointer_convergence.rs index 4e1864aa..197f0b9e 100644 --- a/tests/pointer_convergence.rs +++ b/tests/pointer_convergence.rs @@ -6,6 +6,12 @@ //! These are the property tests behind that claim, plus the two anti-abuse //! properties the merge rule exists to provide: one payment funds one state, //! and no re-signature of a stored state can displace it. +//! +//! ADR-0018 carves out one exception: a final state (counter `u64::MAX`) is +//! replaced by nothing, so two *different* final states are unordered and a +//! node keeps the first one it took. The claim holds for every set with at +//! most one final state in it; with two, the first final state delivered wins, +//! and nothing delivered after it moves a node off it. #![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] @@ -559,11 +565,10 @@ fn an_unknown_version_is_refused_rather_than_accepted_at_its_own_price() { } } -/// At `u64::MAX` no *counter* can out-rank the winner, but a smaller *target* -/// still can. That asymmetry is why migration has to happen before the terminal -/// update rather than on it. +/// At `u64::MAX` nothing out-ranks the held state: no counter is larger, and +/// an equal counter no longer resolves by target (ADR-0018). #[test] -fn a_terminal_counter_cannot_be_out_counted_only_out_targeted() { +fn a_final_state_is_out_ranked_by_nothing() { let terminal = signed(13, u64::MAX, 5, PointerTargetKind::Chunk); assert!(terminal.is_terminal()); assert!(terminal.next_counter().is_err(), "no successor exists"); @@ -575,27 +580,22 @@ fn a_terminal_counter_cannot_be_out_counted_only_out_targeted() { assert!(terminal.replaces(&earlier)); } - // The order does not degenerate there: equal-counter conflicts at the - // maximum still resolve deterministically, so replicas cannot split. + // Nor does another final state, whichever way its target sorts: the first + // one a node took is the one it keeps. let low = signed(13, u64::MAX, 1, PointerTargetKind::Chunk); let high = signed(13, u64::MAX, 2, PointerTargetKind::Chunk); - assert!(low.replaces(&high)); + assert!(!low.replaces(&high)); assert!(!high.replaces(&low)); - assert_eq!(winner(&[&high, &low]).state_id(), low.state_id()); + assert_eq!(winner(&[&high, &low]).state_id(), high.state_id()); assert_eq!(winner(&[&low, &high]).state_id(), low.state_id()); } -/// Why migration must happen *before* the terminal update. -/// -/// At `u64::MAX` the counter can no longer advance, but the pointer is not -/// frozen: the merge order still resolves equal counters by target bytes, and -/// *smaller* target bytes win. So a migration written at the terminal counter -/// can still be displaced — by the owner, or by anyone replaying an older -/// signed record of theirs with a smaller target. The only safe migration is -/// one made while a successor counter still exists, because a strictly larger -/// counter is the one move nothing can answer. +/// A transfer is final on a node's disk: once stored, the former owner cannot +/// take the address back — not with a later counter, which does not exist, and +/// not with a final state whose target sorts first, which the previous rule +/// would have let displace it. #[tokio::test] -async fn migration_must_happen_before_the_terminal_update() { +async fn a_stored_transfer_cannot_be_taken_back() { use ant_node::pointer::store::{PointerStore, PutOutcome}; use ant_protocol::pointer::PointerTarget; use saorsa_pqc::api::sig::ml_dsa_65; @@ -604,62 +604,117 @@ async fn migration_must_happen_before_the_terminal_update() { let store = PointerStore::new(dir.path()).await.expect("store"); let (pk, sk) = ml_dsa_65().generate_keypair_from_seed(&[14u8; 32]); - let sign_at = |counter: u64, target: PointerTarget| { - Pointer::sign(&sk, &pk, counter, target).expect("sign") - }; - - // One update short of the end: a successor counter still exists. - let penultimate = sign_at( - u64::MAX - 1, + let current = Pointer::create( + &sk, + &pk, PointerTarget::new(PointerTargetKind::Chunk, [0x10u8; 32]), - ); - assert!(!penultimate.is_terminal()); + ) + .expect("create"); assert_eq!( - store.put_bytes(&penultimate.to_bytes()).await.expect("put"), + store.put_bytes(¤t.to_bytes()).await.expect("put"), PutOutcome::Changed ); - // The safe migration: spend the last counter. A strictly larger counter - // beats every target, so nothing at u64::MAX - 1 can answer it. - let migration = sign_at( - penultimate.next_counter().expect("successor exists"), - PointerTarget::new(PointerTargetKind::Pointer, [0x80u8; 32]), - ); + // The handover: one final state, pointing at the recipient's pointer. + let recipient = [0x80u8; 32]; + let transfer = current.transfer_to(&sk, recipient).expect("transfer"); + assert_eq!(transfer.transferred_to(), Some(recipient)); assert_eq!( - store.put_bytes(&migration.to_bytes()).await.expect("put"), + store.put_bytes(&transfer.to_bytes()).await.expect("put"), PutOutcome::Changed ); - assert!(migration.is_terminal()); - assert!(migration.next_counter().is_err()); - // Now the danger. The counter is spent, so the only remaining moves are to - // strictly smaller target bytes — and they still win. - let smaller_target = sign_at( + // A former owner grinding a target that sorts first. Under the previous + // rule this displaced the transfer; now it is stale. + let take_back = Pointer::sign( + &sk, + &pk, u64::MAX, - PointerTarget::new(PointerTargetKind::Chunk, [0x01u8; 32]), - ); - assert!( - smaller_target.to_bytes() != migration.to_bytes(), - "a genuinely different state" - ); + PointerTarget::new(PointerTargetKind::Chunk, [0x00u8; 32]), + ) + .expect("sign"); + assert!(take_back.target().to_bytes() < transfer.target().to_bytes()); assert_eq!( - store - .put_bytes(&smaller_target.to_bytes()) - .await - .expect("put"), - PutOutcome::Changed, - "a terminal pointer is NOT frozen: a smaller target still displaces it" + store.put_bytes(&take_back.to_bytes()).await.expect("put"), + PutOutcome::Stale, + "a second final state must not displace the first" ); - // Larger target bytes cannot claw it back: the move is one-way. + // Nor any lower counter, and the held state is the transfer throughout. + for counter in [0u64, 1, u64::MAX - 1] { + let older = Pointer::sign( + &sk, + &pk, + counter, + PointerTarget::new(PointerTargetKind::Chunk, [0x00u8; 32]), + ) + .expect("sign"); + assert_eq!( + store.put_bytes(&older.to_bytes()).await.expect("put"), + PutOutcome::Stale + ); + } assert_eq!( - store.put_bytes(&migration.to_bytes()).await.expect("put"), - PutOutcome::Stale, - "and the displaced migration can never be restored" + store.state_id(&transfer.address()), + Some(transfer.state_id()) ); +} + +proptest! { + #![proptest_config(ProptestConfig::with_cases(16))] - // Which is the whole point: a migration made at the terminal counter is - // not final, so it has to be made earlier, where the counter still answers. - assert!(smaller_target.replaces(&migration)); - assert!(!migration.replaces(&smaller_target)); + /// With one final state among the records, every delivery order still + /// ends on it: rule 0 only leaves two *final* states unordered. + #[test] + fn one_final_state_wins_in_every_delivery_order( + counters in prop::collection::vec(0u64..4, 1..5), + targets in prop::collection::vec(0u8..4, 1..5), + final_target in 0u8..4, + ) { + let mut records: Vec = counters + .iter() + .zip(targets.iter()) + .map(|(counter, target)| signed(1, *counter, *target, PointerTargetKind::Chunk)) + .collect(); + let finalized = signed(1, u64::MAX, final_target, PointerTargetKind::Pointer); + records.push(finalized.clone()); + + let mut order: Vec<&Pointer> = records.iter().collect(); + let mut permutations = 0usize; + permute(&mut order, 0, &mut |candidate| { + prop_assert_eq!(winner(candidate).state_id(), finalized.state_id()); + Ok(()) + }, &mut permutations)?; + prop_assert_eq!(permutations, factorial(records.len())); + } + + /// With two different final states, the one delivered first is what a + /// node keeps, whatever else arrives before, between or after them. + #[test] + fn the_first_final_state_delivered_is_kept( + counters in prop::collection::vec(0u64..4, 1..4), + targets in prop::collection::vec(0u8..4, 1..4), + first_target in 0u8..8, + second_target in 8u8..16, + ) { + let mut records: Vec = counters + .iter() + .zip(targets.iter()) + .map(|(counter, target)| signed(1, *counter, *target, PointerTargetKind::Chunk)) + .collect(); + records.push(signed(1, u64::MAX, first_target, PointerTargetKind::Pointer)); + records.push(signed(1, u64::MAX, second_target, PointerTargetKind::Pointer)); + + let mut order: Vec<&Pointer> = records.iter().collect(); + let mut permutations = 0usize; + permute(&mut order, 0, &mut |candidate| { + let first_final = candidate + .iter() + .find(|record| record.is_terminal()) + .expect("two finals were delivered"); + prop_assert_eq!(winner(candidate).state_id(), first_final.state_id()); + Ok(()) + }, &mut permutations)?; + prop_assert_eq!(permutations, factorial(records.len())); + } } From b8eeb5876ed617d08beb4cde07d49442cdd16af7 Mon Sep 17 00:00:00 2001 From: grumbach Date: Tue, 29 Sep 2026 16:45:58 +0900 Subject: [PATCH 2/8] test(pointer): import what the finality tests use, rather than naming it inline --- src/pointer/service.rs | 13 +++++++------ tests/pointer_convergence.rs | 7 +------ 2 files changed, 8 insertions(+), 12 deletions(-) diff --git a/src/pointer/service.rs b/src/pointer/service.rs index 2578b8e0..effd5683 100644 --- a/src/pointer/service.rs +++ b/src/pointer/service.rs @@ -471,8 +471,9 @@ fn cross_kind_refusal(address: XorName, chunk_present: Result) -> Option

(MlDsaPublicKey, MlDsaSecretKey) { ml_dsa_65().generate_keypair_from_seed(&[seed; 32]) @@ -857,19 +858,19 @@ mod tests { /// which counts how often it was asked. struct StubWitness { conflict: Option, - asked: std::sync::atomic::AtomicUsize, + asked: AtomicUsize, } impl StubWitness { fn new(conflict: Option) -> Arc { Arc::new(Self { conflict, - asked: std::sync::atomic::AtomicUsize::new(0), + asked: AtomicUsize::new(0), }) } fn asked(&self) -> usize { - self.asked.load(std::sync::atomic::Ordering::SeqCst) + self.asked.load(Ordering::SeqCst) } } @@ -878,7 +879,7 @@ mod tests { &'a self, _state: &'a PointerState, ) -> BoxFuture<'a, Option> { - self.asked.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + self.asked.fetch_add(1, Ordering::SeqCst); let conflict = self.conflict.clone(); Box::pin(async move { conflict }) } @@ -887,7 +888,7 @@ mod tests { fn final_state(seed: u8, target_byte: u8) -> Pointer { let (pk, sk) = keypair(seed); let target = PointerTarget::new(PointerTargetKind::Pointer, [target_byte; 32]); - Pointer::sign(&sk, &pk, ant_protocol::pointer::FINAL_COUNTER, target).expect("sign") + Pointer::sign(&sk, &pk, FINAL_COUNTER, target).expect("sign") } #[tokio::test] diff --git a/tests/pointer_convergence.rs b/tests/pointer_convergence.rs index 197f0b9e..6f011c53 100644 --- a/tests/pointer_convergence.rs +++ b/tests/pointer_convergence.rs @@ -17,6 +17,7 @@ use std::collections::BTreeSet; +use ant_node::pointer::store::{PointerStore, PutOutcome}; use ant_protocol::pointer::{Pointer, PointerTarget, PointerTargetKind}; use proptest::prelude::*; use saorsa_pqc::api::sig::{ @@ -444,8 +445,6 @@ fn the_wire_format_is_what_the_adr_says() { /// layer does, which is where the cost would actually have been paid. #[tokio::test] async fn sixty_four_signatures_buy_exactly_one_write() { - use ant_node::pointer::store::{PointerStore, PutOutcome}; - let dir = tempfile::tempdir().expect("tempdir"); let store = PointerStore::new(dir.path()).await.expect("open store"); @@ -596,10 +595,6 @@ fn a_final_state_is_out_ranked_by_nothing() { /// would have let displace it. #[tokio::test] async fn a_stored_transfer_cannot_be_taken_back() { - use ant_node::pointer::store::{PointerStore, PutOutcome}; - use ant_protocol::pointer::PointerTarget; - use saorsa_pqc::api::sig::ml_dsa_65; - let dir = tempfile::tempdir().expect("tempdir"); let store = PointerStore::new(dir.path()).await.expect("store"); From 833c686c9ddb6c0c3e058761d6010e8a8237c358 Mon Sep 17 00:00:00 2001 From: grumbach Date: Tue, 29 Sep 2026 16:56:28 +0900 Subject: [PATCH 3/8] docs(adr): name ADR-0018's pull requests --- docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md b/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md index 642bd7df..1568b810 100644 --- a/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md +++ b/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md @@ -6,7 +6,7 @@ - **Reviewers:** TBD - **Supersedes:** none. Amends ADR-0016's merge rule at the final counter. - **Superseded by:** none -- **Related:** ADR-0016 (pointers), ADR-0005 (repair quorum); WithAutonomi/ant-protocol `feat/pointer-ownership-transfer`, WithAutonomi/ant-client `feat/pointer-ownership-transfer` +- **Related:** ADR-0016 (pointers), ADR-0005 (repair quorum); V2-1354; WithAutonomi/ant-protocol#40, WithAutonomi/ant-node#239, WithAutonomi/ant-client#210 ## Context From 21d308ecc06311161f445f6434109ea0b2a6b974 Mon Sep 17 00:00:00 2001 From: grumbach Date: Tue, 29 Sep 2026 17:04:46 +0900 Subject: [PATCH 4/8] docs(pointer): make the finality check's budget public, as its docs link to it --- src/replication/pointer.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/replication/pointer.rs b/src/replication/pointer.rs index 1fb49aaa..fb23aa8f 100644 --- a/src/replication/pointer.rs +++ b/src/replication/pointer.rs @@ -117,7 +117,7 @@ const MAX_PRUNE_CANDIDATES_PER_PASS: usize = 256; /// to leave that PUT well inside the client's ten-second store timeout. A /// group that has not answered by then is taken to hold nothing that /// conflicts: silence is never a vote, here as anywhere else. -pub(crate) const FINAL_STATE_CHECK_BUDGET: Duration = Duration::from_secs(4); +pub const FINAL_STATE_CHECK_BUDGET: Duration = Duration::from_secs(4); /// One peer's answer about one address: the state it holds there, if any. type StateAnswer = ((PeerId, XorName), Option); From 834feb26a6209390ea3394a4f89e7eeb4f289eb8 Mon Sep 17 00:00:00 2001 From: grumbach Date: Wed, 30 Sep 2026 12:01:31 +0900 Subject: [PATCH 5/8] docs(adr): leave ADR-0016 as written, and say in ADR-0018 what it changes ADR-0016 was edited in place: a backlink, the terminal-counter conclusion, the merge rule and the ownership consequence. It is restored exactly, and ADR-0018 now lists the three statements of ADR-0016 it overrides at the final counter. ADR-0018 also claimed more than the rule gives. The former owner keeps the earlier record and the key, so it can sign a second final state at any time, not only in a race, and any node the first has not reached will take it. What holds is that no node gives up a final state it holds. The claim that flooding final states pays for every round trip was false while a verified payment is cached; the ADR now describes the proven-conflict memory and the bounded, per-address looks that make it true, and the restore of a node's own lost final state without a look. --- docs/adr/ADR-0016-pointers-immutable-owner.md | 17 ++-- ...8-pointer-transfer-by-final-redirection.md | 78 +++++++++++++++---- 2 files changed, 67 insertions(+), 28 deletions(-) diff --git a/docs/adr/ADR-0016-pointers-immutable-owner.md b/docs/adr/ADR-0016-pointers-immutable-owner.md index 8c5648ff..ac1a1909 100644 --- a/docs/adr/ADR-0016-pointers-immutable-owner.md +++ b/docs/adr/ADR-0016-pointers-immutable-owner.md @@ -3,7 +3,7 @@ - **Status:** Proposed - **Date:** 2026-09-18 - **Decision owners:** Anselme (@grumbach) -- **Related:** ADR-0002 (audit), ADR-0004 (commitment-bound pricing), ADR-0008 (per-record pricing), ADR-0009 (audit families), ADR-0011 (capacity-gated discovery), ADR-0014 (file store), ADR-0015 (browser clients), ADR-0018 (a final state is final, and transfer by final redirection — amends the merge rule below) +- **Related:** ADR-0002 (audit), ADR-0004 (commitment-bound pricing), ADR-0008 (per-record pricing), ADR-0009 (audit families), ADR-0011 (capacity-gated discovery), ADR-0014 (file store), ADR-0015 (browser clients) ## Context @@ -85,8 +85,8 @@ rule either would refuse every later update for good, and three such peers would leave no write able to reach its quorum. Under the merge rule each takes the next update, however far ahead of it that is. -An owner who jumps straight to `u64::MAX` limits only themself: the state is -final (ADR-0018), so it is the last one that pointer will ever hold. +An owner who jumps straight to `u64::MAX` limits only themself, and does not +even freeze the pointer: equal counters still resolve by target, below. ### Merge @@ -95,10 +95,6 @@ final (ADR-0018), so it is the last one that pointer will ever hold. 2. smaller target bytes ``` -ADR-0018 puts one rule ahead of these: a state at `u64::MAX` is replaced by -nothing, so two different final states are unordered and a node keeps the -first it took. Below the final counter what follows holds unchanged. - A total order on the states of **one address**. Records of different owners are not comparable and never contend. **Equal state never replaces**: ML-DSA signing is randomized, so one state has unboundedly many valid encodings, and ordering @@ -249,10 +245,9 @@ verified when it was committed and every read verifies it again. self-proving, so one copy settles it; a pointer read has to decide which of several signed states is current, and that answer has to come from more than one peer. -- **The owner key cannot change.** Handover is indirection: point at a new - pointer the recipient owns. Below the final counter the old owner can take - that back, so it is a revocable forwarding; ADR-0018 makes it stick by - signing it at the final counter. +- **Ownership cannot change.** Handover is indirection: point at a new pointer + the recipient owns. The old owner keeps write access forever, so it is a + revocable forwarding state, not a sale. - **Key compromise is permanent.** No rotation, no recovery. - The inlined key costs 1,952 bytes on every read, forever — a deliberate trade for self-contained validation. diff --git a/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md b/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md index 1568b810..c2a76da5 100644 --- a/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md +++ b/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md @@ -38,9 +38,10 @@ to* can. - The address readers use must not change. - No new record type, field, message or storage format; nothing a node has to interpret beyond the counter it already compares. -- Once the network holds a transfer, the former owner must have no move left. -- Forks the owner can still make must be detectable by any reader, and must be - impossible to make once a transfer is established. +- Once a node holds a transfer, no arrival may move it off, the former + owner's included. +- Forks the owner can still make must be detectable by any reader, and must + not be able to take a transfer back from the nodes that hold it. - Below the final counter nothing changes: ADR-0016's convergence argument still holds there. @@ -82,6 +83,20 @@ This lives in `ant_protocol::pointer::PointerState::replaces`, so the node's store, its admission gate, fresh offers, repair and hints all take it without a line of their own. +### What this changes in ADR-0016 + +ADR-0016 is left as written. Where the two disagree, this ADR holds, at the +final counter only: + +- ADR-0016's merge rule gains rule 0 above. +- "An owner who jumps straight to `u64::MAX` ... does not even freeze the + pointer" no longer holds: a state at `u64::MAX` is the last one the pointer + holds on every node that takes it. +- "Ownership cannot change ... a revocable forwarding state, not a sale" no + longer holds for a forwarding signed at the final counter: no node that holds + it gives it up. The owner key still cannot change, and a forwarding below the + final counter is still revocable. + ### Transfer A transfer is the final state whose target is another pointer: @@ -104,8 +119,9 @@ pointer. ### The fork that remains -Only the owner can sign a final state, so only the owner can fork one: sign two -and send them to different nodes before either has heard of the other. Each +Only the owner can sign a final state, so only the owner can fork one. It +keeps the earlier record and its key, so it can sign a second final state at +once or much later, and any node the first has not reached will take it. Each node keeps its first. Nothing local can settle that, and this design does not try. What it guarantees instead: @@ -117,14 +133,30 @@ try. What it guarantees instead: state each holds. A peer claiming a *different* final state is asked for the record, and if it verifies — only the owner could have signed it — the node refuses, answering `Stale` with the state the group proved. A claim alone - refuses nothing, so one dishonest peer cannot block a transfer. The look runs - after payment is verified (an owner flooding its own finals pays for each - round trip it causes), asks only peers that have sent a pointer message, and - is bounded at four seconds inside the client's ten-second store timeout; - silence proves nothing and the write goes ahead. This is what closes the gap - the merge rule leaves: a node that joined the group after the transfer, or - lost its copy, would otherwise take a second final state on the merge rule - alone. + refuses nothing, so one dishonest peer cannot block a transfer. Each peer's + question and fetch run as one pipeline, all at once, so a peer that claims a + rival and then stalls its fetch cannot hide another peer's proof until the + budget runs out. The look runs after payment is verified, asks only peers + that have sent a pointer message, and is bounded at four seconds; silence + proves nothing and the write goes ahead. This is what closes the gap the + merge rule leaves: a node that joined the group after the transfer would + otherwise take a second final state on the merge rule alone. +- **A look costs one round, once.** A payment proof, once verified, is cached, + so replaying one paid final state costs its sender nothing after the first + time. A node therefore remembers each final state a look proved (up to + 16,384 addresses, oldest forgotten first) and refuses a replayed loser from + that memory, ahead of the signature check and without asking anyone. Looks + for one address wait their turn, so a burst of replays costs one round; at + most 64 run at once, and one that cannot start within two seconds answers + the PUT with a retryable error and drops the fresh offer, neither taking nor + refusing the state for good. Two seconds of waiting and four of looking stay + inside the client's ten-second store timeout. Flooding past that needs a + new paid final state per round. +- **A node restores its own final state without looking.** A node that lost + the file of a final state it held is admitted that exact state again and + nothing else (ADR-0016's lost-record rule), so taking it back is a restore, + not a new final state. It asks nobody: a peer on the other side of a fork + would otherwise keep it from restoring its own copy for good. - **Readers see the majority, and see forks.** A client read that meets a final state is settled only once one final state is held by a majority of the close group, and returns that one. If two different final states are seen and @@ -171,8 +203,8 @@ try. What it guarantees instead: changes. - No new record, field, message or storage format. The wire is untouched; the only protocol change is one comparison. -- After the handover lands the former owner has no move at all: no larger - counter exists and no equal one replaces. +- Once the handover lands on a node, the former owner has no move there: no + larger counter exists and no equal one replaces. - A handover can be chained: the recipient can transfer its own pointer on. - ADR-0016's "revocable forwarding" and "migrate before the final counter" caveats are gone: migration *is* the final update. @@ -194,13 +226,18 @@ try. What it guarantees instead: that first round, or one whose group cannot answer in four seconds, takes a second final state on the merge rule alone. Reads still return the majority side; the residual risk is a minority fork that `pointer_finality` reports. +- Under a flood of distinct paid final states a node answers some honest final + PUTs with a retryable error rather than look for them late. The client + retries, as for any refused store. - A final PUT costs its node one state query per capable close-group peer, and a fetch per claimed rival, before it commits. - **Mixed fleets.** A node on ADR-0016's rule still lets a smaller-target final state displace the first. Until a close group's majority runs this rule a former owner can win back the nodes that do not. Reads follow the majority, so the handover holds wherever most of the group has upgraded; recipients should - wait for the upgrade before relying on a transfer. + wait for the upgrade before relying on a transfer. Nothing on the wire tells + the two rules apart: the pointer format version is unchanged, so no node can + refuse to replicate with a peer on the old rule. - Transferred reads take one extra hop, and a chain of transfers one per hop, bounded by the client's resolve depth. @@ -228,7 +265,14 @@ try. What it guarantees instead: - Node, request handler: a final state the group proves already superseded is refused and named, and nothing is written; one nobody contradicts is taken; the group is asked only about a final state the node lacks; two final states - raced to two nodes leave each on its first. + raced to two nodes leave each on its first; a node restores its own lost + final state without asking, and still refuses the rival; a busy look neither + takes nor refuses a final state for good; a proven loser is refused again + before its signature is checked, and without asking. +- Node, the look: a peer that claims a rival and stalls its fetch does not hide + another peer's proof; a claim the served record does not back is no proof; + proven final states refuse others, not themselves, and the oldest is + forgotten first past the cap. - Node, repair: a node holding a final state adopts nothing else; of two final states with quorum, the larger side is adopted in either answer order. - Node, live network: a transfer written to one node reaches the group and no From d8cb7df610c65b37723ce1e24b1f7b56a435f84a Mon Sep 17 00:00:00 2001 From: grumbach Date: Wed, 30 Sep 2026 12:01:31 +0900 Subject: [PATCH 6/8] fix(pointer): make the look before a final state hard to stall, cheap to replay, and safe to restore The finality look took claims as they arrived but fetched each claimed rival in turn. A peer that claimed a rival and then stalled its fetch held the look until its four-second budget ran out, and a look that runs out finds nothing, so one peer could let a second final state in. Each peer's claim and fetch now run as one pipeline, all at once, and the first proof wins. A verified payment is cached, so replaying one paid final state that lost cost the node a signature check and a round of questions every time. A proven conflict is now remembered, up to 16,384 addresses, and a replayed loser is refused before its signature is checked. Looks for one address wait their turn, at most 64 run at once, and one that cannot start within two seconds answers Busy: a PUT gets a retryable error and a fresh offer is dropped, neither taking nor refusing the state for good. A node that lost the file of a final state it held was made to look before restoring it, so a peer on the other side of a fork could keep it from restoring its own copy. The store now reports what it remembers, and restoring a remembered final state asks nobody. ReplicationEngine::with_pointers keeps the signature it has on main; with_pointer_service wires the service both ways. ant-protocol is repinned to the head of its companion change, which only corrects documentation. --- Cargo.lock | 2 +- Cargo.toml | 2 +- src/node.rs | 2 +- src/pointer/mod.rs | 12 +- src/pointer/service.rs | 289 ++++++++++++++++++++++------ src/pointer/store.rs | 10 + src/replication/mod.rs | 43 +++-- src/replication/pointer.rs | 380 +++++++++++++++++++++++++++++++------ tests/e2e/testnet.rs | 2 +- 9 files changed, 605 insertions(+), 137 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index c94a4c29..0c94e357 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -883,7 +883,7 @@ dependencies = [ [[package]] name = "ant-protocol" version = "3.0.0" -source = "git+https://github.com/WithAutonomi/ant-protocol?rev=6f53bb71c62002a312a373281761fe0fdc0fc3aa#6f53bb71c62002a312a373281761fe0fdc0fc3aa" +source = "git+https://github.com/WithAutonomi/ant-protocol?rev=2f2731b9d13d286ab0b9eaf60b316618ccdf29ce#2f2731b9d13d286ab0b9eaf60b316618ccdf29ce" dependencies = [ "blake3", "bytes", diff --git a/Cargo.toml b/Cargo.toml index cc98b10f..11aa670e 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -233,7 +233,7 @@ webrtc-direct = [ # the published 3.0.0, so this is the only entry that has to leave the release # baseline. A rev, not a branch, so the pin is immutable. Drop it once the # pointer PR lands and a release includes it. Versions are bumped at release. -ant-protocol = { git = "https://github.com/WithAutonomi/ant-protocol", rev = "6f53bb71c62002a312a373281761fe0fdc0fc3aa" } +ant-protocol = { git = "https://github.com/WithAutonomi/ant-protocol", rev = "2f2731b9d13d286ab0b9eaf60b316618ccdf29ce" } [profile.release] lto = true diff --git a/src/node.rs b/src/node.rs index ada1065d..ebf9cb91 100644 --- a/src/node.rs +++ b/src/node.rs @@ -275,7 +275,7 @@ impl NodeBuilder { // hands each newly stored paid state to it, and asks it before taking // a final state (ADR-0018). if let Some(service) = protocol.pointer_service() { - engine.with_pointers(service); + engine.with_pointer_service(service); } // ADR-0004: wire the engine's commitment state as the quote generator's diff --git a/src/pointer/mod.rs b/src/pointer/mod.rs index e017dfa7..b859231d 100644 --- a/src/pointer/mod.rs +++ b/src/pointer/mod.rs @@ -10,11 +10,11 @@ //! //! What the address resolves to can still be handed over for good //! (`docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md`): a state at -//! the final counter is replaced by nothing, so an owner who signs one pointing -//! at the new owner's pointer has no move left. The store gets that from the -//! merge rule; the service adds one look at the close group before taking a -//! final state, so a second one cannot land on a node that had not heard of -//! the first. +//! the final counter is replaced by nothing, so once a node holds one pointing +//! at the new owner's pointer, the former owner cannot take it back there. The +//! store gets that from the merge rule; the service adds one look at the close +//! group before taking a final state, so a node that had not heard of the +//! first refuses a second one when a peer proves the first in time. //! //! # What lives here //! @@ -64,5 +64,5 @@ pub use ant_protocol::pointer::{ PointerTarget, PointerTargetKind, DATA_TYPE_POINTER, FINAL_COUNTER, POINTER_BODY_LEN, POINTER_FORMAT_VERSION, POINTER_WIRE_LEN, TARGET_WIRE_LEN, }; -pub use service::{FinalStateWitness, PointerService}; +pub use service::{FinalStateWitness, FinalityCheck, PointerService}; pub use store::{PointerStore, PutOutcome}; diff --git a/src/pointer/service.rs b/src/pointer/service.rs index effd5683..e8173985 100644 --- a/src/pointer/service.rs +++ b/src/pointer/service.rs @@ -17,17 +17,20 @@ //! collision infeasible, not impossible. A node that holds one kind at an //! address refuses the other rather than silently choosing. //! 4. **A look before a final state.** A state at the final counter is -//! replaced by nothing, so a node that takes one can never be corrected. Before -//! taking one it does not already hold, the node asks its close group -//! whether a *different* final state is already held there, and refuses if a -//! peer proves one with the signed record (ADR-0018). That is what keeps a -//! former owner from handing an address over twice, to nodes that had not -//! heard of the first handover yet. +//! replaced by nothing, so a node that takes one can never be corrected. +//! Before taking one it neither holds nor lost, the node asks its close +//! group whether a *different* final state is already held there, and +//! refuses if a peer proves one with the signed record (ADR-0018). That is +//! what keeps a former owner from handing an address over a second time to +//! a node that had not heard of the first, whenever a peer can prove the +//! first in time. A proven loser is remembered and refused again before its +//! signature is checked. //! //! # Order of work //! //! ```text -//! parse → compare with what is held → verify signature → check payment +//! parse → compare with what is held → (final state only) a proven loser? +//! → verify signature → check payment //! → (final state only) ask the close group → commit //! ``` //! @@ -55,7 +58,7 @@ use crate::pointer::store::{Inspected, PointerStore, PutOutcome}; use crate::replication::admission; use crate::replication::pointer::PointerFreshWrite; use crate::storage::{ChunkStore, SELF_CLOSENESS_GATE_WIDTH}; -use ant_protocol::pointer::{Pointer, PointerState, POINTER_WIRE_LEN}; +use ant_protocol::pointer::{PointerState, POINTER_WIRE_LEN}; /// Where a node looks, before it takes a final state, for a different final /// state its close group already holds (ADR-0018). @@ -64,13 +67,32 @@ use ant_protocol::pointer::{Pointer, PointerState, POINTER_WIRE_LEN}; /// built after this service and needs a running P2P node, and because what the /// service decides from the answer is worth testing without one. pub trait FinalStateWitness: Send + Sync { - /// A final state at `state.address`, other than `state`, that a peer in - /// the close group holds — as the signed record, which only the owner can - /// have made — or `None` if none was found. + /// Ask the close group whether a final state other than `state` is + /// already held at `state.address`. + fn check_final<'a>(&'a self, state: &'a PointerState) -> BoxFuture<'a, FinalityCheck>; + + /// A final state at `state.address`, other than `state`, that an earlier + /// check already proved, without asking anyone. /// - /// `None` also when the group could not be asked in time. Silence is not - /// evidence, so it never refuses a write; only a record does. - fn conflicting_final<'a>(&'a self, state: &'a PointerState) -> BoxFuture<'a, Option>; + /// Cheap enough to consult ahead of the signature check, so a final state + /// that lost once is refused again for nothing, however often the same + /// paid record is replayed. + fn proven_conflict(&self, state: &PointerState) -> Option; +} + +/// What a close group said about a final state a node is about to take. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum FinalityCheck { + /// No peer proved a different final state. Also the answer when the group + /// could not be asked in time: silence is not evidence, so it never + /// refuses a write. + Clear, + /// A different final state a peer holds, proven by the signed record, + /// which only the owner could have made. + Conflict(PointerState), + /// Too many checks were running to start this one in time. The state is + /// neither taken nor refused for good, and the sender may try again. + Busy, } /// Handles pointer requests against a [`PointerStore`]. @@ -230,12 +252,8 @@ impl PointerService { // Last, because it costs the close group a round trip, and only for a // paid final state: if a peer proves a different final state is // already held, this one lost the race to it everywhere that matters. - // Named as the state this node knows of, as for any stale arrival. - if let Some(conflict) = self.conflicting_final(&state).await { - return PointerPutResponse::Stale { - address, - state_id: conflict.state_id(), - }; + if let Some(refusal) = self.final_refusal(&state).await { + return refusal; } // Charge the bytes this write will take, and hold the charge until it @@ -360,35 +378,79 @@ impl PointerService { ))); } } + + // A final state an earlier check proved lost to another is refused + // again here, before it buys a signature check or a round trip. + if self.needs_final_check(state) { + let proven = self + .final_witness + .read() + .as_ref() + .and_then(|witness| witness.proven_conflict(state)); + if let Some(conflict) = proven { + return Some(PointerPutResponse::Stale { + address, + state_id: conflict.state_id, + }); + } + } None } - /// The different final state the close group proves it already holds, if - /// `state` is a final state this node would be taking for the first time. + /// Whether `state` needs the close group's word before this node takes + /// it: it is final, and this node neither holds nor remembers a final + /// state at its address. /// - /// Not asked for anything else. A lower counter can be replaced, so taking - /// one blind costs nothing a later state cannot fix; and a node that - /// already holds a final state has had its answer from the merge rule — - /// the same state is unchanged and any other is stale — before this runs. - async fn conflicting_final(&self, state: &PointerState) -> Option { - if !state.is_terminal() - || self + /// Nothing else is asked about. A lower counter can be replaced, so taking + /// one blind costs nothing a later state cannot fix. A node that holds a + /// final state has had its answer from the merge rule before this runs: + /// the same state is unchanged and any other is stale. And a node that + /// lost the file of a final state it held is admitted only that state + /// again ([`PointerStore::admits`]), which restores what it already took + /// rather than taking a new one. Asking there would let any peer on the + /// other side of a fork keep it from restoring its own copy. + fn needs_final_check(&self, state: &PointerState) -> bool { + state.is_terminal() + && !self .store - .state(&state.address) - .is_some_and(|held| held.is_terminal()) - { + .remembered(&state.address) + .is_some_and(|known| known.is_terminal()) + } + + /// The answer to a paid final state the close group proves already lost, + /// or that could not be checked in time; `None` to take it. + async fn final_refusal(&self, state: &PointerState) -> Option { + if !self.needs_final_check(state) { return None; } let witness = self.final_witness.read().as_ref().map(Arc::clone)?; - let conflict = witness.conflicting_final(state).await?; - info!( - "Refusing final pointer state {} at {}: the close group already holds final \ - state {}, so the owner has finalized it before", - hex::encode(state.state_id), - hex::encode(state.address), - hex::encode(conflict.state_id()) - ); - Some(conflict) + match witness.check_final(state).await { + FinalityCheck::Clear => None, + // Named as the state the group proved, as for any stale arrival. + FinalityCheck::Conflict(conflict) => { + info!( + "Refusing final pointer state {} at {}: the close group already holds \ + final state {}, so the owner has finalized it before", + hex::encode(state.state_id), + hex::encode(state.address), + hex::encode(conflict.state_id) + ); + Some(PointerPutResponse::Stale { + address: state.address, + state_id: conflict.state_id, + }) + } + FinalityCheck::Busy => { + debug!( + "Deferring final pointer state {} at {}: too many finality checks running", + hex::encode(state.state_id), + hex::encode(state.address) + ); + Some(PointerPutResponse::Error(ProtocolError::Internal( + "too many finality checks running; try again".to_string(), + ))) + } + } } /// Handle a pointer GET. @@ -857,14 +919,25 @@ mod tests { /// A close group whose answer to the finality question is fixed, and /// which counts how often it was asked. struct StubWitness { - conflict: Option, + answer: FinalityCheck, + proven: Option, asked: AtomicUsize, } impl StubWitness { - fn new(conflict: Option) -> Arc { + fn answering(answer: FinalityCheck) -> Arc { + Arc::new(Self { + answer, + proven: None, + asked: AtomicUsize::new(0), + }) + } + + /// One that already proved `conflict` in an earlier check. + fn having_proven(conflict: PointerState) -> Arc { Arc::new(Self { - conflict, + answer: FinalityCheck::Clear, + proven: Some(conflict), asked: AtomicUsize::new(0), }) } @@ -875,13 +948,15 @@ mod tests { } impl FinalStateWitness for StubWitness { - fn conflicting_final<'a>( - &'a self, - _state: &'a PointerState, - ) -> BoxFuture<'a, Option> { + fn check_final<'a>(&'a self, _state: &'a PointerState) -> BoxFuture<'a, FinalityCheck> { self.asked.fetch_add(1, Ordering::SeqCst); - let conflict = self.conflict.clone(); - Box::pin(async move { conflict }) + let answer = self.answer; + Box::pin(async move { answer }) + } + + fn proven_conflict(&self, state: &PointerState) -> Option { + self.proven + .filter(|proven| proven.state_id != state.state_id) } } @@ -903,7 +978,7 @@ mod tests { let established = final_state(1, 0xAA); let late = final_state(1, 0x01); - let witness = StubWitness::new(Some(established.clone())); + let witness = StubWitness::answering(FinalityCheck::Conflict(established.state())); service.attach_final_state_witness(witness.clone()); match service.handle_put(put(&late)).await { @@ -930,7 +1005,7 @@ mod tests { // Silence is not evidence: a group that proves nothing lets the final // state through on the merge rule. let (service, _dir) = service().await; - let witness = StubWitness::new(None); + let witness = StubWitness::answering(FinalityCheck::Clear); service.attach_final_state_witness(witness.clone()); let transfer = final_state(2, 0x42); @@ -952,7 +1027,7 @@ mod tests { // from the merge rule before any round trip: the same state is // unchanged, and any other is stale. let (service, _dir) = service().await; - let witness = StubWitness::new(None); + let witness = StubWitness::answering(FinalityCheck::Clear); service.attach_final_state_witness(witness.clone()); for counter in [0u64, 1, u64::MAX - 1] { @@ -984,6 +1059,114 @@ mod tests { assert_eq!(witness.asked(), 1, "a held final state answers for itself"); } + #[tokio::test] + async fn a_node_restores_its_own_lost_final_state_without_asking() { + // The node took a final state, then lost the file. The group holds + // the other side of a fork, which a check would prove. Asking would + // keep the node from restoring what it already took, for good; the + // lost state itself is the one thing it is admitted, so nobody is + // asked, and the rival is still refused. + let (service, _dir) = service().await; + let own = final_state(5, 0x10); + let rival = final_state(5, 0x01); + assert!(matches!( + service.handle_put(put(&own)).await, + PointerPutResponse::Success { .. } + )); + std::fs::remove_file(service.store().file_for(&own.address())).expect("remove"); + assert!(service + .store() + .get(&own.address()) + .await + .expect("get") + .is_none()); + + let witness = StubWitness::answering(FinalityCheck::Conflict(rival.state())); + service.attach_final_state_witness(witness.clone()); + assert!(matches!( + service.handle_put(put(&own)).await, + PointerPutResponse::Success { .. } + )); + assert_eq!(witness.asked(), 0, "restoring its own state asks nobody"); + assert_eq!( + service.store().state_id(&own.address()), + Some(own.state_id()) + ); + assert!(matches!( + service.handle_put(put(&rival)).await, + PointerPutResponse::Stale { .. } + )); + } + + #[tokio::test] + async fn a_busy_check_neither_takes_nor_refuses_a_final_state_for_good() { + let (service, _dir) = service().await; + let current = signed(6, 3, 1); + service.handle_put(put(¤t)).await; + let witness = StubWitness::answering(FinalityCheck::Busy); + service.attach_final_state_witness(witness.clone()); + + let transfer = final_state(6, 0x42); + assert!(matches!( + service.handle_put(put(&transfer)).await, + PointerPutResponse::Error(ProtocolError::Internal(_)) + )); + assert_eq!(witness.asked(), 1); + assert_eq!( + service.store().state_id(&transfer.address()), + Some(current.state_id()), + "nothing was written" + ); + assert!( + service.store().admits(&transfer.state()), + "and the state can still be taken later" + ); + } + + #[tokio::test] + async fn a_proven_conflict_is_refused_again_before_any_signature_check() { + // Replaying a paid final state that already lost costs a node + // nothing: the refusal comes from what an earlier check proved, ahead + // of the signature check and without asking the group. A record + // whose signature does not even verify shows the order. + let (unwitnessed, _other) = service().await; + let (service, _dir) = service().await; + let established = final_state(7, 0xAA); + let late = final_state(7, 0x01); + let witness = StubWitness::having_proven(established.state()); + service.attach_final_state_witness(witness.clone()); + + let mut forged = late.to_bytes(); + if let Some(last) = forged.last_mut() { + *last ^= 0xFF; + } + match service + .handle_put(PointerPutRequest::new(Bytes::from(forged.clone()))) + .await + { + PointerPutResponse::Stale { state_id, .. } => { + assert_eq!(state_id, established.state_id()); + } + other => panic!("a proven loser must be refused as stale, got {other:?}"), + } + assert_eq!(witness.asked(), 0, "and nobody is asked again"); + assert!( + matches!( + unwitnessed + .handle_put(PointerPutRequest::new(Bytes::from(forged.clone()))) + .await, + PointerPutResponse::Error(_) + ), + "the record really does fail its signature check" + ); + + // The proven state itself is not refused by its own proof. + assert!(matches!( + service.handle_put(put(&established)).await, + PointerPutResponse::Success { .. } + )); + } + #[tokio::test] async fn two_final_states_raced_to_two_nodes_leave_each_on_its_first() { // The one fork ADR-0018 allows: the owner signs two final states and diff --git a/src/pointer/store.rs b/src/pointer/store.rs index 3628d173..3db6576c 100644 --- a/src/pointer/store.rs +++ b/src/pointer/store.rs @@ -562,6 +562,16 @@ impl PointerStore { .map(|entry| entry.state) } + /// The state this node holds at `address`, or held until it lost the + /// file. + /// + /// Unlike [`Self::state`], a lost record still answers: it is the one + /// state, besides a newer one, that [`Self::admits`] lets restore it. + #[must_use] + pub fn remembered(&self, address: &XorName) -> Option { + self.snapshot(address).map(|entry| entry.state) + } + /// Whether a paid PUT of `state` may be taken. /// /// Any state the merge rule prefers to what is held, whatever its counter. diff --git a/src/replication/mod.rs b/src/replication/mod.rs index 6e877d17..525d8198 100644 --- a/src/replication/mod.rs +++ b/src/replication/mod.rs @@ -2215,17 +2215,21 @@ impl ReplicationEngine { }) } - /// Replicate the pointers `service` stores (ADR-0016), and wire the - /// service to this engine both ways: it hands each newly stored paid - /// state here to be offered on, and asks here before it takes a final - /// state (ADR-0018). + /// Replicate pointers too (ADR-0016): the records in `store`, and the + /// fresh writes the pointer PUT handler sends on `fresh_writes`. + /// + /// The PUT handler is not wired to ask this engine anything before it + /// takes a final state (ADR-0018); [`Self::with_pointer_service`] wires + /// both. /// /// Call before [`Self::start`]. - pub fn with_pointers(&mut self, service: &PointerService) { - let (writes, fresh_writes) = mpsc::unbounded_channel(); - service.attach_fresh_writes(writes); - let replication = Arc::new(pointer::PointerReplication::new( - service.store().clone(), + pub fn with_pointers( + &mut self, + store: PointerStore, + fresh_writes: mpsc::UnboundedReceiver, + ) { + self.pointers = Some(Arc::new(pointer::PointerReplication::new( + store, Arc::clone(&self.storage), Arc::clone(&self.p2p_node), Arc::clone(&self.payment_verifier), @@ -2234,13 +2238,26 @@ impl ReplicationEngine { Arc::clone(&self.send_semaphore), self.shutdown.clone(), self.detached_task_tracker.clone(), - )); - let witness: Arc = replication.clone(); - service.attach_final_state_witness(witness); - self.pointers = Some(replication); + ))); self.pointer_fresh_rx = Some(fresh_writes); } + /// Replicate the pointers `service` stores (ADR-0016), and wire the + /// service to this engine both ways: it hands each newly stored paid + /// state here to be offered on, and asks here before it takes a final + /// state (ADR-0018). + /// + /// Call before [`Self::start`]. + pub fn with_pointer_service(&mut self, service: &PointerService) { + let (writes, fresh_writes) = mpsc::unbounded_channel(); + service.attach_fresh_writes(writes); + self.with_pointers(service.store().clone(), fresh_writes); + if let Some(replication) = &self.pointers { + let witness: Arc = replication.clone(); + service.attach_final_state_witness(witness); + } + } + /// The pointer replication, when enabled. Tests use it to drive rounds. #[cfg(any(test, feature = "test-utils"))] #[must_use] diff --git a/src/replication/pointer.rs b/src/replication/pointer.rs index fb23aa8f..dab098fc 100644 --- a/src/replication/pointer.rs +++ b/src/replication/pointer.rs @@ -30,7 +30,7 @@ //! - **Finality.** Before taking a final state it does not hold, from a client //! or a fresh offer, a node asks the close group whether a different final //! state is already held, and refuses if a peer proves one with the signed -//! record (ADR-0018). See [`PointerReplication::conflicting_final`]. +//! record (ADR-0018). See [`PointerReplication::check_final`]. //! //! Requests go only to peers that have sent a pointer message themselves (see //! [`PointerReplication::is_capable`]). A peer built before pointers cannot @@ -39,7 +39,8 @@ //! neighbour-sync round pushes hints, empty or not, so capability is learned //! within a cycle. -use std::collections::{HashMap, HashSet}; +use std::collections::{HashMap, HashSet, VecDeque}; +use std::future::Future; use std::sync::Arc; use std::time::{Duration, Instant}; @@ -62,7 +63,7 @@ use crate::payment::{ MIN_PAYMENT_PROOF_SIZE_BYTES, }; use crate::pointer::store::{Inspected, PointerStore}; -use crate::pointer::FinalStateWitness; +use crate::pointer::{FinalStateWitness, FinalityCheck}; use crate::replication::admission; use crate::replication::commitment_state::ResponderCommitmentState; use crate::replication::config::{ @@ -119,6 +120,26 @@ const MAX_PRUNE_CANDIDATES_PER_PASS: usize = 256; /// conflicts: silence is never a vote, here as anywhere else. pub const FINAL_STATE_CHECK_BUDGET: Duration = Duration::from_secs(4); +/// How long a finality check waits for its turn before it gives up and +/// answers [`FinalityCheck::Busy`]. +/// +/// With [`FINAL_STATE_CHECK_BUDGET`] after it, a check still ends well inside +/// a client's ten-second store timeout. +pub const FINAL_STATE_CHECK_WAIT: Duration = Duration::from_secs(2); + +/// How many finality checks run at once. +/// +/// Each runs on its own share of the address space. Two checks for one +/// address never run together, so a burst of replays of one final state +/// costs the group one round of questions. +pub const FINAL_STATE_CHECK_STRIPES: usize = 64; + +/// How many proven conflicts a node remembers, oldest forgotten first. +/// +/// A conflict is a final state the close group proved, so it never goes +/// stale; the cap only bounds memory, at about 200 bytes an entry. +const MAX_PROVEN_FINALS: usize = 16_384; + /// One peer's answer about one address: the state it holds there, if any. type StateAnswer = ((PeerId, XorName), Option); @@ -246,10 +267,47 @@ pub struct PointerReplication { pending: Mutex>, /// When each held address was first seen continuously out of range. out_of_range: Mutex>, + /// Final states the close group proved, by address (ADR-0018). + proven_finals: Mutex, + /// One turn per share of the address space for a finality check. + final_check_turns: Vec>, shutdown: CancellationToken, tracker: TaskTracker, } +/// Final states a finality check proved the close group holds, remembered so +/// that a paid final state which lost to one is refused again without asking +/// anyone. +#[derive(Default)] +struct ProvenFinals { + by_address: HashMap, + /// Addresses in the order they were first proven, oldest first. + order: VecDeque, +} + +impl ProvenFinals { + /// The proven final state at `state.address`, if it is not `state`. + fn conflict(&self, state: &PointerState) -> Option { + self.by_address + .get(&state.address) + .filter(|proven| proven.state_id != state.state_id) + .copied() + } + + /// Remember `proven`, forgetting the oldest entries past the cap. + fn remember(&mut self, proven: PointerState) { + if self.by_address.insert(proven.address, proven).is_none() { + self.order.push_back(proven.address); + } + while self.by_address.len() > MAX_PROVEN_FINALS { + let Some(oldest) = self.order.pop_front() else { + break; + }; + self.by_address.remove(&oldest); + } + } +} + impl PointerReplication { /// Pointer replication over `store`, sharing the engine's resources. #[allow(clippy::too_many_arguments)] @@ -277,6 +335,10 @@ impl PointerReplication { capable: Mutex::new(HashSet::new()), pending: Mutex::new(HashMap::new()), out_of_range: Mutex::new(HashMap::new()), + proven_finals: Mutex::new(ProvenFinals::default()), + final_check_turns: (0..FINAL_STATE_CHECK_STRIPES) + .map(|_| tokio::sync::Mutex::new(())) + .collect(), shutdown, tracker, } @@ -955,6 +1017,16 @@ impl PointerReplication { { return; } + // A final state already proven to have lost is dropped before it + // buys a signature check. + let needs_final_check = state.is_terminal() + && !self + .store + .remembered(&state.address) + .is_some_and(|known| known.is_terminal()); + if needs_final_check && self.proven_finals.lock().conflict(&state).is_some() { + return; + } let record = match self.store.verify(parsed).await { Ok(record) => record, Err(e) => { @@ -979,25 +1051,13 @@ impl PointerReplication { ); return; } - // A final state held nowhere here yet is looked for in the group - // first, exactly as a client PUT of one is: the peer offering it may - // simply be the side of a race this node has not heard the other - // side of. - let holds_final = self - .store - .state(&state.address) - .is_some_and(|held| held.is_terminal()); - if state.is_terminal() && !holds_final { - if let Some(conflict) = self.conflicting_final(&state).await { - info!( - "Refusing fresh final pointer state {} at {} from {source}: the close \ - group already holds final state {}", - hex::encode(state.state_id), - hex::encode(state.address), - hex::encode(conflict.state_id()) - ); - return; - } + // A final state this node neither holds nor remembers is looked for + // in the group first, exactly as a client PUT of one is: the peer + // offering it may simply be the side of a race this node has not + // heard the other side of. A node that lost its own final state is + // admitted only that state, so restoring it asks nobody. + if needs_final_check && !self.offered_final_is_clear(&source, &state).await { + return; } match self .store_verified(record, Some(self.config.paid_list_close_group_size)) @@ -1019,14 +1079,42 @@ impl PointerReplication { // Finality // ----------------------------------------------------------------------- - /// A final state at `state.address`, other than `state`, that a peer in - /// the close group holds and proves by serving the signed record. + /// Whether a freshly offered final state may be taken: no peer proved a + /// different one, and the look was not too busy to run. + async fn offered_final_is_clear(&self, source: &PeerId, state: &PointerState) -> bool { + match self.check_final(state).await { + FinalityCheck::Clear => true, + FinalityCheck::Conflict(conflict) => { + info!( + "Refusing fresh final pointer state {} at {} from {source}: the close \ + group already holds final state {}", + hex::encode(state.state_id), + hex::encode(state.address), + hex::encode(conflict.state_id) + ); + false + } + FinalityCheck::Busy => { + debug!( + "Dropping fresh final pointer state {} at {} from {source}: too many \ + finality checks running", + hex::encode(state.state_id), + hex::encode(state.address) + ); + false + } + } + } + + /// Ask the close group whether a final state at `state.address`, other + /// than `state`, is already held, and have a peer that says so prove it + /// by serving the signed record. /// /// Asked before this node takes a final state it does not hold. A final /// state is replaced by nothing, so a node that took a second one would /// hold it for good; asking first is what keeps a former owner from - /// finalizing an address again on nodes that had not yet heard it was - /// final — one that joined the group since, or lost its copy. + /// finalizing an address again on a node that had not yet heard it was + /// final, such as one that joined the group since. /// /// A peer's word is not enough to refuse: a state summary is a claim /// anyone can make, and one dishonest peer could otherwise block every @@ -1034,30 +1122,56 @@ impl PointerReplication { /// final state, so one that verifies is the owner's own proof that it /// finalized the pointer before. /// - /// Bounded by [`FINAL_STATE_CHECK_BUDGET`], and only peers that have sent - /// a pointer message are asked. A group that cannot be asked in time finds - /// nothing, and the write goes ahead on the merge rule alone: a race is a - /// fork the client detects, not one this can prevent. - pub async fn conflicting_final(&self, state: &PointerState) -> Option { + /// A proof is remembered, so the same loser is refused again without a + /// question. Checks for one address wait their turn, so a burst of replays + /// costs one round, and at most [`FINAL_STATE_CHECK_STRIPES`] run at once; + /// one that cannot start within [`FINAL_STATE_CHECK_WAIT`] answers `Busy`. + /// Once started it is bounded by [`FINAL_STATE_CHECK_BUDGET`], and only + /// peers that have sent a pointer message are asked. A group that cannot + /// be asked in time finds nothing, and the write goes ahead on the merge + /// rule alone: a race is a fork the client detects, not one this can + /// prevent. + pub async fn check_final(&self, state: &PointerState) -> FinalityCheck { if !state.is_terminal() { - return None; + return FinalityCheck::Clear; + } + let proven = self.proven_finals.lock().conflict(state); + if let Some(conflict) = proven { + return FinalityCheck::Conflict(conflict); + } + let stripe = usize::from(state.address.first().copied().unwrap_or_default()) + % FINAL_STATE_CHECK_STRIPES; + let Some(turn) = self.final_check_turns.get(stripe) else { + return FinalityCheck::Busy; + }; + let Ok(_turn) = tokio::time::timeout(FINAL_STATE_CHECK_WAIT, turn.lock()).await else { + return FinalityCheck::Busy; + }; + // A check that held the turn may have proven it meanwhile. + let proven = self.proven_finals.lock().conflict(state); + if let Some(conflict) = proven { + return FinalityCheck::Conflict(conflict); } - tokio::time::timeout(FINAL_STATE_CHECK_BUDGET, self.find_conflicting_final(state)) + match tokio::time::timeout(FINAL_STATE_CHECK_BUDGET, self.find_conflicting_final(state)) .await - .unwrap_or_else(|_| { + { + Ok(Some(record)) => { + let conflict = record.state(); + self.proven_finals.lock().remember(conflict); + FinalityCheck::Conflict(conflict) + } + Ok(None) => FinalityCheck::Clear, + Err(_) => { debug!( "Close group of pointer {} did not answer the finality check in time", hex::encode(state.address) ); - None - }) + FinalityCheck::Clear + } + } } - /// The body of [`Self::conflicting_final`], without its time bound. - /// - /// Answers are taken as they arrive, and a claim is checked the moment it - /// is made, so one slow peer cannot hide a quick one's proof behind the - /// budget. + /// The body of [`Self::check_final`], without its bounds. async fn find_conflicting_final(&self, state: &PointerState) -> Option { let self_id = *self.p2p.peer_id(); let address = state.address; @@ -1070,23 +1184,13 @@ impl PointerReplication { .map(|node| node.peer_id) .filter(|peer| *peer != self_id && self.is_capable(peer)) .collect(); - - let mut claims = stream::iter(peers) - .map(|peer| async move { - let held = self.ask_state(&peer, &address).await?; - (held.is_terminal() && held.state_id != state.state_id).then_some(peer) - }) - .buffer_unordered(VERIFICATION_CONCURRENCY); - while let Some(claim) = claims.next().await { - let Some(peer) = claim else { continue }; - let Some(record) = self.fetch_record(&peer, &address).await else { - continue; - }; - if record.is_terminal() && record.state_id() != state.state_id { - return Some(record); - } - } - None + first_proven_final( + peers, + state, + |peer| async move { self.ask_state(&peer, &address).await }, + |peer| async move { self.fetch_record(&peer, &address).await }, + ) + .await } /// Ask one peer which state it holds at `address`. @@ -1318,16 +1422,62 @@ impl PointerReplication { } impl FinalStateWitness for PointerReplication { - fn conflicting_final<'a>(&'a self, state: &'a PointerState) -> BoxFuture<'a, Option> { - Box::pin(Self::conflicting_final(self, state)) + fn check_final<'a>(&'a self, state: &'a PointerState) -> BoxFuture<'a, FinalityCheck> { + Box::pin(Self::check_final(self, state)) + } + + fn proven_conflict(&self, state: &PointerState) -> Option { + self.proven_finals.lock().conflict(state) + } +} + +/// The first final state other than `state` that one of `peers` claims with +/// `ask` and then proves with `fetch`. +/// +/// Each peer's question and fetch run as one pipeline, all of them at once, +/// and the first proof wins. A peer that claims a rival and then stalls its +/// fetch holds up nobody else's proof: were the fetches made one at a time +/// after the claims, that one peer would hold the check until its budget ran +/// out, and a finality check that runs out finds nothing. +async fn first_proven_final( + peers: Vec, + state: &PointerState, + ask: Ask, + fetch: Fetch, +) -> Option +where + Ask: Fn(PeerId) -> AskFut + Sync, + AskFut: Future> + Send, + Fetch: Fn(PeerId) -> FetchFut + Sync, + FetchFut: Future> + Send, +{ + let rival = |held: &PointerState| held.is_terminal() && held.state_id != state.state_id; + let (ask, fetch, rival) = (&ask, &fetch, &rival); + let width = peers.len().max(1); + let mut proofs = stream::iter(peers) + .map(|peer| async move { + let held = ask(peer).await?; + if !rival(&held) { + return None; + } + let record = fetch(peer).await?; + rival(&record.state()).then_some(record) + }) + .buffer_unordered(width); + while let Some(proof) = proofs.next().await { + if proof.is_some() { + return proof; + } } + None } #[cfg(test)] #[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] mod tests { use super::*; - use ant_protocol::pointer::{PointerTarget, PointerTargetKind}; + use ant_protocol::pointer::{PointerTarget, PointerTargetKind, FINAL_COUNTER}; + use saorsa_pqc::api::sig::ml_dsa_65; fn state(counter: u64, target: u8, id: u8) -> PointerState { PointerState { @@ -1516,4 +1666,112 @@ mod tests { fn a_group_nobody_can_be_asked_in_is_undecided() { assert_eq!(evaluate(None, 0, &[], 0), Verdict::Undecided); } + + fn final_record(seed: u8, target: u8) -> Pointer { + let (pk, sk) = ml_dsa_65().generate_keypair_from_seed(&[seed; 32]); + let target = PointerTarget::new(PointerTargetKind::Pointer, [target; 32]); + Pointer::sign(&sk, &pk, FINAL_COUNTER, target).expect("sign") + } + + #[tokio::test(start_paused = true)] + async fn a_claimant_that_stalls_its_fetch_does_not_hide_a_rival_another_peer_proves() { + // One peer claims a rival final state first and then never serves + // it; an honest peer claims the same rival a moment later and serves + // it at once. The honest proof must land inside the budget. + let taking = final_record(9, 0x01); + let rival = final_record(9, 0xAA); + let (staller, honest) = (peer(1), peer(2)); + let found = tokio::time::timeout( + FINAL_STATE_CHECK_BUDGET, + first_proven_final( + vec![staller, honest], + &taking.state(), + |asked| { + let claim = rival.state(); + async move { + if asked == honest { + tokio::time::sleep(Duration::from_millis(100)).await; + } + Some(claim) + } + }, + |asked| { + let record = rival.clone(); + async move { + if asked == staller { + std::future::pending::<()>().await; + } + Some(record) + } + }, + ), + ) + .await; + assert_eq!( + found.ok().flatten().map(|record| record.state_id()), + Some(rival.state_id()), + "the honest peer's proof was hidden behind the stalled fetch" + ); + } + + #[tokio::test] + async fn only_a_different_final_state_that_is_served_is_a_proof() { + let taking = final_record(10, 0x01); + let rival = final_record(10, 0xAA); + let (same, unbacked, open) = (peer(1), peer(2), peer(3)); + let found = first_proven_final( + vec![same, unbacked, open], + &taking.state(), + |asked| { + let claim = if asked == same { + taking.state() + } else if asked == unbacked { + rival.state() + } else { + state(3, 1, 1) + }; + async move { Some(claim) } + }, + // The peer claiming a rival serves the state being taken instead. + |_| { + let record = taking.clone(); + async move { Some(record) } + }, + ) + .await; + assert!( + found.is_none(), + "a claim the record does not back is no proof" + ); + } + + #[test] + fn a_proven_final_state_refuses_others_and_forgets_the_oldest_past_the_cap() { + let mut proven = ProvenFinals::default(); + let held = final_record(11, 0xAA).state(); + let other = final_record(11, 0x01).state(); + proven.remember(held); + assert_eq!(proven.conflict(&other), Some(held)); + assert_eq!( + proven.conflict(&held), + None, + "a proof is no conflict with itself" + ); + + for i in 0..MAX_PROVEN_FINALS { + let mut filler = held; + filler.address = [0; 32]; + if let Some(slot) = filler.address.get_mut(..8) { + slot.copy_from_slice(&(i as u64 + 1).to_be_bytes()); + } + proven.remember(filler); + } + assert_eq!(proven.by_address.len(), MAX_PROVEN_FINALS); + assert_eq!(proven.order.len(), MAX_PROVEN_FINALS); + assert_eq!( + proven.conflict(&other), + None, + "the oldest proof is forgotten first" + ); + } } diff --git a/tests/e2e/testnet.rs b/tests/e2e/testnet.rs index 64d7c5e3..a20c225f 100644 --- a/tests/e2e/testnet.rs +++ b/tests/e2e/testnet.rs @@ -1375,7 +1375,7 @@ impl TestNetwork { // which the service also asks before a final state // (ADR-0018). if let Some(service) = protocol.pointer_service() { - engine.with_pointers(service); + engine.with_pointer_service(service); } let dht_events = p2p.dht_manager().subscribe_events(); engine.start(dht_events); From c5166eb1c0c490ed1ec1e0ea527071cda311d0a7 Mon Sep 17 00:00:00 2001 From: grumbach Date: Wed, 30 Sep 2026 12:31:46 +0900 Subject: [PATCH 7/8] chore(deps): pin ant-protocol at the head of its companion change The companion change added documentation only: where a final state is permanent, and that conflicts below the final counter are ordered. --- Cargo.lock | 2 +- Cargo.toml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 0c94e357..1cbcd9bf 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -883,7 +883,7 @@ dependencies = [ [[package]] name = "ant-protocol" version = "3.0.0" -source = "git+https://github.com/WithAutonomi/ant-protocol?rev=2f2731b9d13d286ab0b9eaf60b316618ccdf29ce#2f2731b9d13d286ab0b9eaf60b316618ccdf29ce" +source = "git+https://github.com/WithAutonomi/ant-protocol?rev=b4d2c7c7ebb243a97afc7b6e6e36de98c4ff18de#b4d2c7c7ebb243a97afc7b6e6e36de98c4ff18de" dependencies = [ "blake3", "bytes", diff --git a/Cargo.toml b/Cargo.toml index 11aa670e..77e39058 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -233,7 +233,7 @@ webrtc-direct = [ # the published 3.0.0, so this is the only entry that has to leave the release # baseline. A rev, not a branch, so the pin is immutable. Drop it once the # pointer PR lands and a release includes it. Versions are bumped at release. -ant-protocol = { git = "https://github.com/WithAutonomi/ant-protocol", rev = "2f2731b9d13d286ab0b9eaf60b316618ccdf29ce" } +ant-protocol = { git = "https://github.com/WithAutonomi/ant-protocol", rev = "b4d2c7c7ebb243a97afc7b6e6e36de98c4ff18de" } [profile.release] lto = true From e08f9ded384fd1abef15dea1a6e31266dd49ae25 Mon Sep 17 00:00:00 2001 From: grumbach Date: Wed, 30 Sep 2026 13:15:57 +0900 Subject: [PATCH 8/8] fix(pointer): never let a second proof replace the first, and share one look among queued replays A proven final state was stored one per address, so a later look that proved the other side of a fork replaced the first proof, and the state the first look had disproved could then be taken if the group went quiet. Up to two proven states are now kept per address, and a second never replaces the first; two different ones refuse every final state there. Looks for one address took turns but did not share answers, so replays queued behind a clear look each asked the group again once it finished. A clear look now answers the same state for ten seconds without asking, on the runtime's clock, so a burst of replays costs one round as ADR-0018 says. The bounds live in their own type so they are tested without a network: eight queued checks make one look, and a proof answers every replay of the loser. --- ...8-pointer-transfer-by-final-redirection.md | 31 +- src/replication/pointer.rs | 266 ++++++++++++++---- 2 files changed, 226 insertions(+), 71 deletions(-) diff --git a/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md b/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md index c2a76da5..2fa086af 100644 --- a/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md +++ b/docs/adr/ADR-0018-pointer-transfer-by-final-redirection.md @@ -143,15 +143,19 @@ try. What it guarantees instead: otherwise take a second final state on the merge rule alone. - **A look costs one round, once.** A payment proof, once verified, is cached, so replaying one paid final state costs its sender nothing after the first - time. A node therefore remembers each final state a look proved (up to - 16,384 addresses, oldest forgotten first) and refuses a replayed loser from - that memory, ahead of the signature check and without asking anyone. Looks - for one address wait their turn, so a burst of replays costs one round; at - most 64 run at once, and one that cannot start within two seconds answers - the PUT with a retryable error and drops the fresh offer, neither taking nor - refusing the state for good. Two seconds of waiting and four of looking stay - inside the client's ten-second store timeout. Flooding past that needs a - new paid final state per round. + time. A node therefore remembers each final state a look proved, up to two + per address and 16,384 addresses, oldest forgotten first, and never lets a + later proof replace an earlier one; it refuses a replayed loser from that + memory, ahead of the signature check and without asking anyone. A look that + found nothing answers the same state again for ten seconds without asking. + Looks for one address wait their turn, so replays queued behind a look find + its answer, and a burst costs one round; at most 64 run at once, and one + that cannot start within two seconds answers the PUT with an error and + drops the fresh offer, neither taking nor refusing the state for good, so + the write can be tried again. + Two seconds of waiting and four of looking stay inside the client's + ten-second store timeout. Flooding past that needs a new paid final state + per round. - **A node restores its own final state without looking.** A node that lost the file of a final state it held is admitted that exact state again and nothing else (ADR-0016's lost-record rule), so taking it back is a restore, @@ -227,8 +231,9 @@ try. What it guarantees instead: second final state on the merge rule alone. Reads still return the majority side; the residual risk is a minority fork that `pointer_finality` reports. - Under a flood of distinct paid final states a node answers some honest final - PUTs with a retryable error rather than look for them late. The client - retries, as for any refused store. + PUTs with an error rather than look for them late. A client retries one such + refusal as a shortfall; several make the write fail, and the owner tries + again. - A final PUT costs its node one state query per capable close-group peer, and a fetch per claimed rival, before it commits. - **Mixed fleets.** A node on ADR-0016's rule still lets a smaller-target final @@ -272,7 +277,9 @@ try. What it guarantees instead: - Node, the look: a peer that claims a rival and stalls its fetch does not hide another peer's proof; a claim the served record does not back is no proof; proven final states refuse others, not themselves, and the oldest is - forgotten first past the cap. + forgotten first past the cap; a second proof at an address never replaces + the first; replays queued behind a look reuse its answer, and a proof + answers every replay of the loser. - Node, repair: a node holding a final state adopts nothing else; of two final states with quorum, the larger side is adopted in either answer order. - Node, live network: a transfer written to one node reaches the group and no diff --git a/src/replication/pointer.rs b/src/replication/pointer.rs index dab098fc..9d7e7c09 100644 --- a/src/replication/pointer.rs +++ b/src/replication/pointer.rs @@ -53,6 +53,7 @@ use saorsa_core::identity::PeerId; use saorsa_core::{P2PNode, TrustEvent}; use tokio::sync::{mpsc, RwLock, Semaphore}; use tokio::task::JoinHandle; +use tokio::time::Instant as TokioInstant; use tokio_util::sync::CancellationToken; use tokio_util::task::TaskTracker; @@ -134,12 +135,28 @@ pub const FINAL_STATE_CHECK_WAIT: Duration = Duration::from_secs(2); /// costs the group one round of questions. pub const FINAL_STATE_CHECK_STRIPES: usize = 64; -/// How many proven conflicts a node remembers, oldest forgotten first. +/// How many addresses a node remembers proven final states for, oldest +/// forgotten first. /// -/// A conflict is a final state the close group proved, so it never goes -/// stale; the cap only bounds memory, at about 200 bytes an entry. +/// A proof is a final state the close group was shown to hold, so it never +/// goes stale; the cap only bounds memory, at most two states of about 200 +/// bytes each per address. const MAX_PROVEN_FINALS: usize = 16_384; +/// Proven final states remembered per address. Two different ones are enough +/// to refuse every final state there: each conflicts with the other, and any +/// third with both. +const MAX_PROOFS_PER_ADDRESS: usize = 2; + +/// How long a look that found no conflict answers again for the same final +/// state without asking anyone: long enough to cover replays queued behind +/// it, short against how fast a close group changes. +pub const FINAL_STATE_CLEAR_REUSE: Duration = Duration::from_secs(10); + +/// Most clear looks remembered for reuse. Each is at most +/// [`FINAL_STATE_CLEAR_REUSE`] old, so this only bounds a burst. +const MAX_CLEAR_LOOKS: usize = 4096; + /// One peer's answer about one address: the state it holds there, if any. type StateAnswer = ((PeerId, XorName), Option); @@ -267,10 +284,8 @@ pub struct PointerReplication { pending: Mutex>, /// When each held address was first seen continuously out of range. out_of_range: Mutex>, - /// Final states the close group proved, by address (ADR-0018). - proven_finals: Mutex, - /// One turn per share of the address space for a finality check. - final_check_turns: Vec>, + /// The bounds on looks before a final state (ADR-0018). + finality: FinalityLooks, shutdown: CancellationToken, tracker: TaskTracker, } @@ -280,23 +295,33 @@ pub struct PointerReplication { /// anyone. #[derive(Default)] struct ProvenFinals { - by_address: HashMap, + /// Up to [`MAX_PROOFS_PER_ADDRESS`] different proven states per address. + by_address: HashMap>, /// Addresses in the order they were first proven, oldest first. order: VecDeque, } impl ProvenFinals { - /// The proven final state at `state.address`, if it is not `state`. + /// A proven final state at `state.address` other than `state`. fn conflict(&self, state: &PointerState) -> Option { self.by_address - .get(&state.address) - .filter(|proven| proven.state_id != state.state_id) + .get(&state.address)? + .iter() + .find(|proven| proven.state_id != state.state_id) .copied() } - /// Remember `proven`, forgetting the oldest entries past the cap. + /// Remember `proven` beside what is already proven at its address, never + /// in place of it: a proof replaced would let the state it disproved in. + /// Forgets the oldest addresses past the cap. fn remember(&mut self, proven: PointerState) { - if self.by_address.insert(proven.address, proven).is_none() { + if let Some(states) = self.by_address.get_mut(&proven.address) { + let known = states.iter().any(|state| state.state_id == proven.state_id); + if !known && states.len() < MAX_PROOFS_PER_ADDRESS { + states.push(proven); + } + } else { + self.by_address.insert(proven.address, vec![proven]); self.order.push_back(proven.address); } while self.by_address.len() > MAX_PROVEN_FINALS { @@ -308,6 +333,107 @@ impl ProvenFinals { } } +/// The bounds on looks before a final state (ADR-0018), apart from the +/// network so they can be tested without one. +/// +/// A look's answer is kept either way: a proof for good, a clear look briefly. +/// Looks for one address take turns, so replays queued behind a look find its +/// answer rather than asking again; at most [`FINAL_STATE_CHECK_STRIPES`] +/// run at once. +struct FinalityLooks { + proven: Mutex, + /// When each final state was last looked for and found clear, on the + /// runtime's clock, which the look's own time bounds use as well. + clear: Mutex>, + /// One turn per share of the address space. + turns: Vec>, +} + +impl FinalityLooks { + fn new() -> Self { + Self { + proven: Mutex::new(ProvenFinals::default()), + clear: Mutex::new(HashMap::new()), + turns: (0..FINAL_STATE_CHECK_STRIPES) + .map(|_| tokio::sync::Mutex::new(())) + .collect(), + } + } + + /// A proven final state that conflicts with `state`, asking nobody. + fn proven_conflict(&self, state: &PointerState) -> Option { + self.proven.lock().conflict(state) + } + + /// Whether a look found `state` clear recently enough to answer again. + fn recently_clear(&self, state: &PointerState, now: TokioInstant) -> bool { + self.clear + .lock() + .get(&(state.address, state.state_id)) + .is_some_and(|at| now.saturating_duration_since(*at) < FINAL_STATE_CLEAR_REUSE) + } + + fn note_clear(&self, state: &PointerState, now: TokioInstant) { + let mut clear = self.clear.lock(); + if clear.len() >= MAX_CLEAR_LOOKS { + clear.retain(|_, at| now.saturating_duration_since(*at) < FINAL_STATE_CLEAR_REUSE); + } + if clear.len() < MAX_CLEAR_LOOKS { + clear.insert((state.address, state.state_id), now); + } + } + + /// Answer whether `state` may be taken, running `look` only when neither + /// a proof nor a recent clear look answers it already. + async fn check(&self, state: &PointerState, look: Look) -> FinalityCheck + where + Look: FnOnce() -> Fut + Send, + Fut: Future> + Send, + { + if !state.is_terminal() { + return FinalityCheck::Clear; + } + if let Some(conflict) = self.proven_conflict(state) { + return FinalityCheck::Conflict(conflict); + } + let stripe = usize::from(state.address.first().copied().unwrap_or_default()) + % FINAL_STATE_CHECK_STRIPES; + let Some(turn) = self.turns.get(stripe) else { + return FinalityCheck::Busy; + }; + let Ok(_turn) = tokio::time::timeout(FINAL_STATE_CHECK_WAIT, turn.lock()).await else { + return FinalityCheck::Busy; + }; + // A look that held the turn may have answered this meanwhile. + if let Some(conflict) = self.proven_conflict(state) { + return FinalityCheck::Conflict(conflict); + } + if self.recently_clear(state, TokioInstant::now()) { + return FinalityCheck::Clear; + } + match tokio::time::timeout(FINAL_STATE_CHECK_BUDGET, look()).await { + Ok(Some(record)) => { + let conflict = record.state(); + self.proven.lock().remember(conflict); + FinalityCheck::Conflict(conflict) + } + Ok(None) => { + self.note_clear(state, TokioInstant::now()); + FinalityCheck::Clear + } + // Silence proves nothing, and is not kept for reuse either: the + // next look may be answered. + Err(_) => { + debug!( + "Close group of pointer {} did not answer the finality check in time", + hex::encode(state.address) + ); + FinalityCheck::Clear + } + } + } +} + impl PointerReplication { /// Pointer replication over `store`, sharing the engine's resources. #[allow(clippy::too_many_arguments)] @@ -335,10 +461,7 @@ impl PointerReplication { capable: Mutex::new(HashSet::new()), pending: Mutex::new(HashMap::new()), out_of_range: Mutex::new(HashMap::new()), - proven_finals: Mutex::new(ProvenFinals::default()), - final_check_turns: (0..FINAL_STATE_CHECK_STRIPES) - .map(|_| tokio::sync::Mutex::new(())) - .collect(), + finality: FinalityLooks::new(), shutdown, tracker, } @@ -1024,7 +1147,7 @@ impl PointerReplication { .store .remembered(&state.address) .is_some_and(|known| known.is_terminal()); - if needs_final_check && self.proven_finals.lock().conflict(&state).is_some() { + if needs_final_check && self.finality.proven_conflict(&state).is_some() { return; } let record = match self.store.verify(parsed).await { @@ -1122,53 +1245,21 @@ impl PointerReplication { /// final state, so one that verifies is the owner's own proof that it /// finalized the pointer before. /// - /// A proof is remembered, so the same loser is refused again without a - /// question. Checks for one address wait their turn, so a burst of replays - /// costs one round, and at most [`FINAL_STATE_CHECK_STRIPES`] run at once; - /// one that cannot start within [`FINAL_STATE_CHECK_WAIT`] answers `Busy`. + /// A proof is remembered, beside any other for the address, so the same + /// loser is refused again without a question; a clear look is reused for + /// the same state for [`FINAL_STATE_CLEAR_REUSE`]. Checks for one address + /// wait their turn, so a burst of replays queued behind a look costs that + /// one round, and at most [`FINAL_STATE_CHECK_STRIPES`] run at once; one + /// that cannot start within [`FINAL_STATE_CHECK_WAIT`] answers `Busy`. /// Once started it is bounded by [`FINAL_STATE_CHECK_BUDGET`], and only /// peers that have sent a pointer message are asked. A group that cannot /// be asked in time finds nothing, and the write goes ahead on the merge /// rule alone: a race is a fork the client detects, not one this can /// prevent. pub async fn check_final(&self, state: &PointerState) -> FinalityCheck { - if !state.is_terminal() { - return FinalityCheck::Clear; - } - let proven = self.proven_finals.lock().conflict(state); - if let Some(conflict) = proven { - return FinalityCheck::Conflict(conflict); - } - let stripe = usize::from(state.address.first().copied().unwrap_or_default()) - % FINAL_STATE_CHECK_STRIPES; - let Some(turn) = self.final_check_turns.get(stripe) else { - return FinalityCheck::Busy; - }; - let Ok(_turn) = tokio::time::timeout(FINAL_STATE_CHECK_WAIT, turn.lock()).await else { - return FinalityCheck::Busy; - }; - // A check that held the turn may have proven it meanwhile. - let proven = self.proven_finals.lock().conflict(state); - if let Some(conflict) = proven { - return FinalityCheck::Conflict(conflict); - } - match tokio::time::timeout(FINAL_STATE_CHECK_BUDGET, self.find_conflicting_final(state)) + self.finality + .check(state, || self.find_conflicting_final(state)) .await - { - Ok(Some(record)) => { - let conflict = record.state(); - self.proven_finals.lock().remember(conflict); - FinalityCheck::Conflict(conflict) - } - Ok(None) => FinalityCheck::Clear, - Err(_) => { - debug!( - "Close group of pointer {} did not answer the finality check in time", - hex::encode(state.address) - ); - FinalityCheck::Clear - } - } } /// The body of [`Self::check_final`], without its bounds. @@ -1427,7 +1518,7 @@ impl FinalStateWitness for PointerReplication { } fn proven_conflict(&self, state: &PointerState) -> Option { - self.proven_finals.lock().conflict(state) + self.finality.proven_conflict(state) } } @@ -1478,6 +1569,7 @@ mod tests { use super::*; use ant_protocol::pointer::{PointerTarget, PointerTargetKind, FINAL_COUNTER}; use saorsa_pqc::api::sig::ml_dsa_65; + use std::sync::atomic::{AtomicUsize, Ordering}; fn state(counter: u64, target: u8, id: u8) -> PointerState { PointerState { @@ -1774,4 +1866,60 @@ mod tests { "the oldest proof is forgotten first" ); } + + #[test] + fn a_second_proof_never_replaces_the_first() { + // Both sides of a fork proven at one address: each is refused by the + // other, so a later look that proved the other side cannot let back + // in the state the first look disproved. + let mut proven = ProvenFinals::default(); + let a = final_record(12, 0xAA).state(); + let b = final_record(12, 0xBB).state(); + let c = final_record(12, 0xCC).state(); + proven.remember(a); + proven.remember(b); + proven.remember(c); + assert_eq!(proven.conflict(&b), Some(a)); + assert_eq!(proven.conflict(&a), Some(b)); + assert!(proven.conflict(&c).is_some(), "a third is refused by both"); + assert_eq!( + proven.by_address.get(&a.address).map(Vec::len), + Some(MAX_PROOFS_PER_ADDRESS) + ); + } + + #[tokio::test(start_paused = true)] + async fn replays_queued_behind_a_look_reuse_its_answer() { + let looks = FinalityLooks::new(); + let taking = final_record(13, 0x01); + let asked = AtomicUsize::new(0); + let look = || { + asked.fetch_add(1, Ordering::SeqCst); + async { + tokio::time::sleep(Duration::from_millis(500)).await; + None + } + }; + let state = taking.state(); + let checks = (0..8).map(|_| looks.check(&state, look)); + let answers = futures::future::join_all(checks).await; + assert!(answers.iter().all(|answer| *answer == FinalityCheck::Clear)); + assert_eq!(asked.load(Ordering::SeqCst), 1, "eight replays, one look"); + + // Past the reuse window the group is asked again. + tokio::time::sleep(FINAL_STATE_CLEAR_REUSE).await; + assert_eq!(looks.check(&state, look).await, FinalityCheck::Clear); + assert_eq!(asked.load(Ordering::SeqCst), 2); + + // A proof answers every replay of the loser without a look. + let loser = final_record(13, 0x02).state(); + let rival = final_record(13, 0xAA); + let proved = looks.check(&loser, || async { Some(rival.clone()) }).await; + assert_eq!(proved, FinalityCheck::Conflict(rival.state())); + assert_eq!( + looks.check(&loser, look).await, + FinalityCheck::Conflict(rival.state()) + ); + assert_eq!(asked.load(Ordering::SeqCst), 2); + } }