From 3b9efc214cd4516b2e0866036c74a53781d66ddd Mon Sep 17 00:00:00 2001 From: Morgan Pretty Date: Tue, 15 Sep 2026 14:04:52 +1000 Subject: [PATCH 1/6] Fetch the user's profile from the whole swarm at once An ordinary poll asks one swarm member, once every poll interval, and the storage server does not promise a config has reached every member of a swarm. Measured against mainnet through a client bridge: an account created minutes earlier, restored on a second device, took 66 seconds for its display name to come back -- several members sampled one at a time before one of them had it. Both mobile clients work around exactly this and say so in their own comments. `fetch_user_profile` asks every member concurrently. It is deliberately asymmetric about what it concludes from what they say: having the config is not a majority property, so one member holding it is the whole answer and the others not having it yet is the condition being routed around rather than evidence against it. An answer carrying the config seconded by a second member ends the fetch early; an empty answer never does, however many members give it; and when everyone has answered, any member that had it still wins. Concluding there is no config is the expensive direction and waits for the whole swarm. Every member is asked the same question -- no `last_hash` -- which is both what makes the answers comparable and what stops two members agreeing on a subset while a third, further along, is still answering. Two pieces of the poll path are now shared rather than duplicated. `decode_retrieved` takes no node and touches no database, deliberately: the retrieve cursor is kept per (namespace, node), and a decoder that knew about nodes is the shape in which one member's cursor gets written from another's answer -- silently, because each cursor is individually plausible. `_record_swarm_cursor` takes the node as a parameter for the same reason. Measured after: 17 seconds, against 66 before. --- include/session/core.hpp | 47 ++++ src/core.cpp | 500 ++++++++++++++++++++++++++++++--------- 2 files changed, 432 insertions(+), 115 deletions(-) diff --git a/include/session/core.hpp b/include/session/core.hpp index 27a03468..a2ceeb65 100644 --- a/include/session/core.hpp +++ b/include/session/core.hpp @@ -348,6 +348,28 @@ class Core { std::string body, int round); + // Records where `node`'s next retrieve of `ns` resumes from. The node is a parameter rather + // than something read from the surroundings because the cursor is kept per (namespace, node): + // writing one member's cursor from another's answer is individually plausible and shows up + // much later as a namespace that re-fetches for ever or one that skips messages. + void _record_swarm_cursor( + const network::ed25519_pubkey& node_pubkey, + config::Namespace ns, + std::span messages); + + // The fan-out behind fetch_user_profile(). Defined in core.cpp: nothing outside it needs the + // shape, and one of them holds a decoded response. + struct ProfileAnswer; + struct ProfileFanOut; + + std::vector _profile_retrieve_body(); + void _handle_profile_response( + ProfileFanOut& state, + const network::service_node& node, + std::optional body, + bool timed_out); + void _settle_profile_fetch(ProfileFanOut& state, ProfileAnswer* taken); + // Decrypts and dispatches one-to-one messages from Namespace::Default. void _handle_direct_messages(std::span messages); @@ -534,6 +556,31 @@ class Core { /// already required to tolerate seeing a message twice). void set_poll_interval(std::chrono::milliseconds interval); + /// Fetches the account's own UserProfile config from the whole swarm at once, rather than + /// waiting for the next poll to ask one member. + /// + /// **For the moment an account arrives on a device that has never had it** -- a restore, where + /// nothing local can answer "what is this account called" and the app has a person waiting in + /// front of a progress indicator. An ordinary poll asks a single member, once every + /// `set_poll_interval`, and the storage server does not promise a config has reached every + /// member of a swarm: measured against mainnet, an account created minutes earlier took over a + /// minute of polling to come back, which is several members sampled one at a time before one of + /// them had it. Both mobile clients work around exactly this, and say so in their own + /// comments. + /// + /// Asks every member concurrently, and is deliberately asymmetric about what it will conclude + /// from what they say. **Having the config is not a majority property**: one member holding it + /// is the whole answer, and the others not having it yet is the condition being routed around + /// rather than evidence against it. So an answer carrying the config, seconded by a second + /// member, ends the fetch early; an empty answer never does, however many members give it, and + /// when everyone has answered any member that had the config still wins. Concluding that there + /// is no config is the expensive direction, and it waits for the whole swarm. + /// + /// `done` is called exactly once, on Core's loop, with whether a config was merged. Safe to + /// call without waiting on it: what it finds is merged into `configs` like anything a poll + /// brings in, and the retrieve cursor it records means the next poll carries on from there. + void fetch_user_profile(std::function done); + /// Encrypt and send a direct message to the given recipient. /// /// Returns a unique message_id that will later be reported via the message_send_status diff --git a/src/core.cpp b/src/core.cpp index f12560d3..30924f99 100644 --- a/src/core.cpp +++ b/src/core.cpp @@ -24,6 +24,7 @@ #include #include #include +#include #include #include "core/swarm_request.hpp" @@ -159,6 +160,124 @@ static constexpr std::array POLL_NAMESPACES = { // fewer; this exists so that a node whose `more` never goes false cannot poll indefinitely. static constexpr int POLL_MAX_ROUNDS = 20; +/// One namespace's slice of a batch retrieve response, decoded but not yet handled. +/// +/// The decoded bytes are kept beside the messages because `SwarmMessage::data` spans into them: +/// the messages are only usable while this object is. +struct retrieved_namespace { + std::vector> data; + std::vector messages; + /// The node says it is holding more past what it returned. A retrieve is capped by the + /// storage server, so this is how a namespace reports that one round did not exhaust it. + bool more = false; +}; + +namespace { + +/// Decodes one result of a batch retrieve, or nothing where the node did not answer it. +/// +/// **Nothing and an empty answer are different**, and the difference is load-bearing: a namespace +/// that failed or answered malformedly is not reported to its handler at all, while one that +/// answered with nothing is -- "we asked and there is nothing" is an answer, and some handlers act +/// on it. +/// +/// Takes no node and touches no database, deliberately. Which node an answer came from matters to +/// whoever called: the retrieve cursor is kept per node and has to stay with the one that produced +/// it, and a decoder that knew about nodes is the shape in which one node's cursor gets written +/// from another node's response -- silently, because each node's cursor is individually plausible. +std::optional decode_retrieved(const nlohmann::json& res, int16_t ns_val) { + auto code_it = res.find("code"); + if (code_it == res.end() || code_it->get() != 200) { + log::warning(cat, "Retrieve of namespace {} failed: {}", ns_val, res.dump()); + return std::nullopt; + } + auto body_it = res.find("body"); + if (body_it == res.end()) + return std::nullopt; + auto msgs_it = body_it->find("messages"); + if (msgs_it == body_it->end() || !msgs_it->is_array()) + return std::nullopt; + + retrieved_namespace got; + if (auto m = body_it->find("more"); m != body_it->end() && m->is_boolean()) + got.more = m->get(); + + log::debug(cat, "Retrieved {} message(s) from namespace {}", msgs_it->size(), ns_val); + + for (const auto& msg : *msgs_it) { + auto data_it = msg.find("data"); + if (data_it == msg.end() || !data_it->is_string()) + continue; + auto& decoded = got.data.emplace_back(); + auto b64 = data_it->get(); + decoded.reserve(oxenc::from_base64_size(b64.size())); + oxenc::from_base64(b64.begin(), b64.end(), std::back_inserter(decoded)); + + SwarmMessage swarm_msg; + swarm_msg.data = {decoded.data(), decoded.size()}; + + if (auto h = msg.find("hash"); h != msg.end() && h->is_string()) + swarm_msg.hash = h->get(); + + if (auto t = msg.find("timestamp"); t != msg.end() && t->is_number_integer()) + swarm_msg.timestamp = from_epoch_ms(t->get()); + + if (auto e = msg.find("expiry"); e != msg.end() && e->is_number_integer()) + swarm_msg.expiry = from_epoch_ms(e->get()); + + got.messages.push_back(std::move(swarm_msg)); + } + + // A node claiming more while returning nothing cannot be continued: there is no new hash to + // move the cursor to, so another round would ask the same question and get the same answer. + // Reported as finished instead, or a handler waiting on `is_final` would wait for one that + // never comes. + got.more = got.more && !got.messages.empty(); + return got; +} + +} // namespace + +/// One distinct answer to the profile fetch, and how many members gave it. +struct Core::ProfileAnswer { + size_t agreed = 0; + /// The member whose answer this is -- the first to give it. Kept because the retrieve cursor + /// is written against the node that produced the messages and no other. + network::ed25519_pubkey node; + retrieved_namespace answer; +}; + +/// What a fan-out has heard so far. +/// +/// Shared between every in-flight request and touched only on Core's loop, which is what makes a +/// plain count safe here: the responses arrive on the network's threads and are marshalled across +/// before any of this is read. +struct Core::ProfileFanOut { + std::function done; + size_t outstanding = 0; + bool settled = false; + /// Keyed by the hashes the answer carried, so agreement is a count rather than a comparison of + /// every answer against every other. + std::unordered_map heard; +}; + +namespace { + +/// What two members have to agree on before their answer is taken. +/// +/// The hashes, sorted: the storage server returns messages in its own order, and two members +/// holding the same config are not obliged to list it identically. +std::string answer_digest(const retrieved_namespace& got) { + std::vector hashes; + hashes.reserve(got.messages.size()); + for (const auto& m : got.messages) + hashes.emplace_back(m.hash); + std::ranges::sort(hashes); + return "{}"_format(fmt::join(hashes, "\n")); +} + +} // namespace + void Core::_poll() { // Non-owning: the Network is ours alone, and callbacks below must not keep it alive -- doing so // could make the loop thread the last owner and run ~Network there. @@ -180,6 +299,198 @@ void Core::_poll() { }); } +/// How many members have to give the same answer before it is taken without waiting for the rest. +/// +/// Two, which is the smallest number that is more than one member's word -- and it gates *having* +/// the config, not the wait. **Absence never settles this fetch early**, however many members +/// report it: presence is not a majority property. The case this whole thing exists for is a +/// config that has reached one member of a swarm and not the others, so counting "I do not have it" +/// against the one member that does would be the original fault with more steps in it. +static constexpr size_t PROFILE_FETCH_QUORUM = 2; + +void Core::fetch_user_profile(std::function done) { + // Non-owning: the Network is ours alone, and callbacks below must not keep it alive -- doing so + // could make the loop thread the last owner and run ~Network there. + auto* net = _network.get(); + if (!net || !globals.have_account()) { + if (done) + done(false); + return; + } + + auto state = std::make_shared(); + state->done = std::move(done); + + net->get_swarm(globals.pubkey_x25519(), false, [this, net, state](auto, auto swarm) { + _loop.call([this, net, state, swarm = std::move(swarm)] { + if (state->settled) + return; + if (swarm.empty()) { + // Not an error worth failing loudly over: the ordinary poll still runs, and this + // was only ever an attempt to get there sooner. A swarm of one is still worth + // asking -- one member holding the config is the whole answer. + log::warning(cat, "Cannot fan out a profile fetch: no swarm members available"); + _settle_profile_fetch(*state, nullptr); + return; + } + + auto body = _profile_retrieve_body(); + state->outstanding = swarm.size(); + + for (const auto& node : swarm) { + net->send_request( + swarm_request(node, globals.pubkey_x25519(), "batch", body), + [this, state, node]( + bool success, + bool timeout, + int16_t /*status_code*/, + std::vector> /*headers*/, + std::optional body) { + // Marshalled rather than handled here: these arrive on the network's + // threads, several at once by design, and everything they touch -- + // the tally, the config merge, the cursor -- belongs to Core's loop. + _loop.call([this, + state, + node, + success, + timeout, + body = std::move(body)]() mutable { + _handle_profile_response( + *state, + node, + (success && body) ? std::move(body) : std::nullopt, + timeout); + }); + }); + } + }); + }); +} + +/// The body of a profile retrieve: one namespace, and no cursor. +/// +/// **No `last_hash`, deliberately.** A cursor says "everything after what I already have", and the +/// point of this fetch is to get the config from a member that may never have been asked before -- +/// resuming from another member's position would ask the wrong question. The same body goes to +/// every member, which is also what makes their answers comparable. +std::vector Core::_profile_retrieve_body() { + auto now_ms = epoch_ms(clock_now_ms()); + auto ns_val = static_cast(config::Namespace::UserProfile); + + nlohmann::json params = { + {"pubkey", globals.session_id_hex()}, + {"namespace", ns_val}, + }; + + if (retrieve_requires_auth(ns_val)) { + auto seed = globals.account_seed(); + auto to_sign = ns_signature_value("retrieve", ns_val, now_ms); + params["pubkey_ed25519"] = globals.pubkey_ed25519().hex(); + params["timestamp"] = now_ms; + params["signature"] = "{:b}"_format(ed25519::sign(seed.ed25519_secret(), to_span(to_sign))); + } + + nlohmann::json requests = nlohmann::json::array(); + requests.push_back({{"method", "retrieve"}, {"params", std::move(params)}}); + return to_vector(nlohmann::json{{"requests", std::move(requests)}}.dump()); +} + +/// One member's answer, tallied against the others. +void Core::_handle_profile_response( + ProfileFanOut& state, + const network::service_node& node, + std::optional body, + bool timed_out) { + if (state.settled) + return; + + state.outstanding -= 1; + + if (body) { + try { + auto json = nlohmann::json::parse(*body); + auto it = json.find("results"); + if (it != json.end() && it->is_array() && !it->empty()) { + if (auto got = decode_retrieved( + (*it)[0], static_cast(config::Namespace::UserProfile))) { + auto& entry = state.heard[answer_digest(*got)]; + if (entry.agreed == 0) { + entry.node = node.remote_pubkey; + entry.answer = std::move(*got); + } + entry.agreed += 1; + + // Seconded, and carrying something: nothing better is going to arrive, so the + // remaining members are not waited on. An empty answer never takes this exit + // however many members give it -- see PROFILE_FETCH_QUORUM. + if (!entry.answer.messages.empty() && entry.agreed >= PROFILE_FETCH_QUORUM) { + _settle_profile_fetch(state, &entry); + return; + } + } + } + } catch (const std::exception& e) { + log::warning(cat, "Failed to parse profile fetch response: {}", e.what()); + } + } else { + log::warning( + cat, + "Profile fetch from {} failed: {}", + node.remote_pubkey.hex(), + timed_out ? "timed out" : "request failed"); + } + + if (state.outstanding > 0) + return; + + // Everyone has answered and no answer was seconded. **Any member that had the config still + // wins**, because the alternative is not a safer answer but no answer at all -- and the + // ordinary poll, which is what runs instead, believes a single member without asking anyone + // else. Members reporting nothing are not counted against it at all: they are the condition + // being routed around, not evidence. + // + // Among answers that do carry something, the most-agreed wins, and among equals the one + // carrying the most messages -- a member holding part of a config is further along than one + // holding less of it. + ProfileAnswer* best = nullptr; + for (auto& [digest, entry] : state.heard) { + if (entry.answer.messages.empty()) + continue; + if (!best || entry.agreed > best->agreed || + (entry.agreed == best->agreed && + entry.answer.messages.size() > best->answer.messages.size())) + best = &entry; + } + _settle_profile_fetch(state, best); +} + +/// Ends the fan-out, merging `taken` if there is one. Called exactly once per fetch. +void Core::_settle_profile_fetch(ProfileFanOut& state, ProfileAnswer* taken) { + if (state.settled) + return; + state.settled = true; + + bool found = taken && !taken->answer.messages.empty(); + + if (taken && !taken->answer.messages.empty()) { + auto configs_held = configs.batch(); + // Final, even where the member said it was holding more. This is a one-shot fetch rather + // than a poll that continues: the cursor recorded below is what lets the ordinary poll pick + // up anything left, and a handler told to wait for a final that never comes would wait for + // ever. + receive_messages(taken->answer.messages, config::Namespace::UserProfile, true); + _record_swarm_cursor(taken->node, config::Namespace::UserProfile, taken->answer.messages); + } + + log::info( + cat, + "Profile fetch settled: {}", + found ? "config merged" : "nothing held by the swarm"); + + if (state.done) + state.done(found); +} + void Core::_send_poll( network::Network* net, network::service_node node, @@ -292,6 +603,74 @@ SELECT h.hash FROM swarm_hashes h JOIN swarm_nodes n ON n.id = h.node }); } +/// Records where this node's next retrieve of `ns` should resume from. +/// +/// **The node is a parameter and not something read from the surroundings, deliberately.** The +/// cursor is kept per (namespace, node) -- the swarm filters a retrieve on a hash that particular +/// member still holds -- so writing one node's cursor from another node's answer is a silent fault: +/// each cursor is individually plausible, and what it costs shows up much later as a namespace that +/// re-fetches for ever or one that skips past messages, depending which way it landed. +/// +/// Called only once the batch has been handled: the swarm filters on last_hash, so advancing past +/// messages that threw would drop them permanently. Handling and then dying before this point +/// re-delivers the batch instead, which is why message handlers must tolerate seeing a message +/// twice. +void Core::_record_swarm_cursor( + const network::ed25519_pubkey& node_pubkey, + config::Namespace ns, + std::span messages) { + if (messages.empty()) + return; + + auto ns_val = static_cast(ns); + auto conn = db.conn(); + + // Every hash goes in, not just the ones that produced something we kept: the cursor is a + // position in what this node returned, so leaving out what we ignored would park it behind + // those and fetch them again on every poll. Insertion order is the order the node returned + // them, which is what `id DESC` reads back. + conn.prepared_exec( + "INSERT INTO swarm_nodes (pubkey) VALUES (?) ON CONFLICT DO NOTHING", node_pubkey); + auto node_id = + conn.prepared_get("SELECT id FROM swarm_nodes WHERE pubkey = ?", node_pubkey); + + for (const auto& m : messages) { + if (m.hash.empty()) + continue; + conn.prepared_exec( + R"( +INSERT INTO swarm_hashes (namespace, node, hash, expiry) VALUES (?, ?, ?, ?) +ON CONFLICT(namespace, node, hash) DO UPDATE SET expiry = max(expiry, excluded.expiry) +)", + ns_val, + node_id, + m.hash, + m.expiry.time_since_epoch().count() > 0 ? std::optional{epoch_ms(m.expiry)} + : std::nullopt); + } + + // An expired hash is not a cursor: the node no longer holds the message to measure from. + conn.prepared_exec( + "DELETE FROM swarm_hashes WHERE expiry IS NOT NULL AND expiry <= ?", + epoch_ms(clock_now_ms())); + + // And a cap on top of that, because expiry alone bounds this at every message in the retention + // window. Only the newest entry is ever read; the rest exist solely to walk back past hashes + // deleted from the swarm, so keeping more than a run of deletions could plausibly cover buys + // nothing but disk. + conn.prepared_exec( + R"( +DELETE FROM swarm_hashes + WHERE namespace = ?1 AND node = ?2 + AND id NOT IN (SELECT id FROM swarm_hashes + WHERE namespace = ?1 AND node = ?2 + ORDER BY id DESC LIMIT ?3) +)", + ns_val, + node_id, + SWARM_HASH_HISTORY); +} + void Core::_handle_poll_response( network::service_node node, std::vector namespaces, @@ -316,129 +695,20 @@ void Core::_handle_poll_response( auto configs_held = configs.batch(); auto& results = *it; - auto conn = db.conn(); for (size_t i = 0; i < namespaces.size() && i < results.size(); ++i) { auto ns = namespaces[i]; auto ns_val = static_cast(ns); - const auto& res = results[i]; - auto code_it = res.find("code"); - if (code_it == res.end() || code_it->get() != 200) { - log::warning(cat, "Retrieve of namespace {} failed: {}", ns_val, res.dump()); - continue; - } - auto body_it = res.find("body"); - if (body_it == res.end()) - continue; - auto msgs_it = body_it->find("messages"); - if (msgs_it == body_it->end() || !msgs_it->is_array()) + auto got = decode_retrieved(results[i], ns_val); + if (!got) continue; - // A retrieve is capped, so this says whether the node is holding more past what it - // returned. Everything above `continue`s instead, which is the distinction that - // matters: a namespace that failed or answered malformedly is not reported to its - // handler at all, while one that answered with nothing is -- "we asked and there is - // nothing" is an answer, and some handlers act on it. - bool more = false; - if (auto m = body_it->find("more"); m != body_it->end() && m->is_boolean()) - more = m->get(); - - log::debug(cat, "Retrieved {} message(s) from namespace {}", msgs_it->size(), ns_val); - - // Decode each message; keep the decoded bytes alive until after - // receive_messages() returns, since SwarmMessage::data spans - // into them. - std::vector> messages_data; - std::vector swarm_messages; - - for (const auto& msg : *msgs_it) { - auto data_it = msg.find("data"); - if (data_it == msg.end() || !data_it->is_string()) - continue; - auto& decoded = messages_data.emplace_back(); - auto b64 = data_it->get(); - decoded.reserve(oxenc::from_base64_size(b64.size())); - oxenc::from_base64(b64.begin(), b64.end(), std::back_inserter(decoded)); - - SwarmMessage swarm_msg; - swarm_msg.data = {decoded.data(), decoded.size()}; - - if (auto h = msg.find("hash"); h != msg.end() && h->is_string()) - swarm_msg.hash = h->get(); - - if (auto t = msg.find("timestamp"); t != msg.end() && t->is_number_integer()) - swarm_msg.timestamp = from_epoch_ms(t->get()); - - if (auto e = msg.find("expiry"); e != msg.end() && e->is_number_integer()) - swarm_msg.expiry = from_epoch_ms(e->get()); - - swarm_messages.push_back(std::move(swarm_msg)); - } - - // A node claiming more while returning nothing cannot be continued: there is no new - // hash to move the cursor to, so another round would ask the same question and get the - // same answer. Treat the namespace as finished instead, or a handler waiting on - // `is_final` would wait for one that never comes. - more = more && !swarm_messages.empty(); - - receive_messages(swarm_messages, ns, !more); - if (more) + const auto& swarm_messages = got->messages; + receive_messages(swarm_messages, ns, !got->more); + if (got->more) unfinished.push_back(ns); - if (!swarm_messages.empty()) { - // Only advance the cursor once the batch has been handled: the swarm filters on - // last_hash, so advancing past messages that threw would drop them permanently. - // Handling then dying before this point re-delivers the batch instead, so message - // handlers must tolerate seeing a message twice. - // - // Every hash goes in, not just the ones that produced something we kept: the cursor - // is a position in what this node returned, so leaving out what we ignored would - // park it behind those and fetch them again on every poll. Insertion order is the - // order the node returned them, which is what `id DESC` reads back. - conn.prepared_exec( - "INSERT INTO swarm_nodes (pubkey) VALUES (?) ON CONFLICT DO NOTHING", - sn_pubkey); - auto node_id = conn.prepared_get( - "SELECT id FROM swarm_nodes WHERE pubkey = ?", sn_pubkey); - - for (const auto& m : swarm_messages) { - if (m.hash.empty()) - continue; - conn.prepared_exec( - R"( -INSERT INTO swarm_hashes (namespace, node, hash, expiry) VALUES (?, ?, ?, ?) -ON CONFLICT(namespace, node, hash) DO UPDATE SET expiry = max(expiry, excluded.expiry) -)", - ns_val, - node_id, - m.hash, - m.expiry.time_since_epoch().count() > 0 - ? std::optional{epoch_ms(m.expiry)} - : std::nullopt); - } - - // An expired hash is not a cursor: the node no longer holds the message to measure - // from. - conn.prepared_exec( - "DELETE FROM swarm_hashes WHERE expiry IS NOT NULL AND expiry <= ?", - epoch_ms(clock_now_ms())); - - // And a cap on top of that, because expiry alone bounds this at every message in - // the retention window. Only the newest entry is ever read; the rest exist solely - // to walk back past hashes deleted from the swarm, so keeping more than a run of - // deletions could plausibly cover buys nothing but disk. - conn.prepared_exec( - R"( -DELETE FROM swarm_hashes - WHERE namespace = ?1 AND node = ?2 - AND id NOT IN (SELECT id FROM swarm_hashes - WHERE namespace = ?1 AND node = ?2 - ORDER BY id DESC LIMIT ?3) -)", - ns_val, - node_id, - SWARM_HASH_HISTORY); - } + _record_swarm_cursor(sn_pubkey, ns, swarm_messages); } } catch (const std::exception& e) { log::warning(cat, "Failed to parse poll response: {}", e.what()); From 1fbaf2c2473c6c933382e6a68e4ccf6f6c3cf71f Mon Sep 17 00:00:00 2001 From: Morgan Pretty Date: Tue, 15 Sep 2026 14:57:22 +1000 Subject: [PATCH 2/6] Take the first member that has the config, without seconding Waiting for a second member to return the same answer bought protection against something that repairs itself. A config is merged rather than assigned, so a stale one taken from a member behind its swarm is corrected by the next ordinary poll; the seconding only delayed the answer, and it delayed it in front of somebody watching a progress indicator. Measured on mainnet, same account, same method: 66s through the ordinary poll, 17s with seconding, 5s without it. The cost of the confirmation was most of the remaining wait. The asymmetry that matters stays: an empty answer still never settles the fetch, however many members give it, and only every member having answered concludes there is nothing to find. --- include/session/core.hpp | 13 ++++-- src/core.cpp | 86 ++++++++-------------------------------- 2 files changed, 25 insertions(+), 74 deletions(-) diff --git a/include/session/core.hpp b/include/session/core.hpp index a2ceeb65..891dbb92 100644 --- a/include/session/core.hpp +++ b/include/session/core.hpp @@ -571,10 +571,15 @@ class Core { /// Asks every member concurrently, and is deliberately asymmetric about what it will conclude /// from what they say. **Having the config is not a majority property**: one member holding it /// is the whole answer, and the others not having it yet is the condition being routed around - /// rather than evidence against it. So an answer carrying the config, seconded by a second - /// member, ends the fetch early; an empty answer never does, however many members give it, and - /// when everyone has answered any member that had the config still wins. Concluding that there - /// is no config is the expensive direction, and it waits for the whole swarm. + /// rather than evidence against it. So the first answer carrying a config ends the fetch, and + /// an empty answer never does however many members give it -- only every member having answered + /// concludes that there is nothing to find. + /// + /// **No agreement is required of the answer that wins**, on purpose. A config is merged rather + /// than assigned, so a stale one taken from a member behind its swarm is corrected by the next + /// ordinary poll rather than stuck; waiting for a second member to say the same thing buys + /// protection against something that already repairs itself, and costs a round trip in front of + /// somebody watching a progress indicator. /// /// `done` is called exactly once, on Core's loop, with whether a config was merged. Safe to /// call without waiting on it: what it finds is merged into `configs` like anything a poll diff --git a/src/core.cpp b/src/core.cpp index 30924f99..7e3cdd9d 100644 --- a/src/core.cpp +++ b/src/core.cpp @@ -24,7 +24,6 @@ #include #include #include -#include #include #include "core/swarm_request.hpp" @@ -238,16 +237,16 @@ std::optional decode_retrieved(const nlohmann::json& res, i } // namespace -/// One distinct answer to the profile fetch, and how many members gave it. +/// The answer a fetch settled on, and the member that gave it. +/// +/// The member is kept with it because the retrieve cursor is written against the node that +/// produced the messages and no other. struct Core::ProfileAnswer { - size_t agreed = 0; - /// The member whose answer this is -- the first to give it. Kept because the retrieve cursor - /// is written against the node that produced the messages and no other. network::ed25519_pubkey node; retrieved_namespace answer; }; -/// What a fan-out has heard so far. +/// How much of a fan-out is still outstanding. /// /// Shared between every in-flight request and touched only on Core's loop, which is what makes a /// plain count safe here: the responses arrive on the network's threads and are marshalled across @@ -256,28 +255,8 @@ struct Core::ProfileFanOut { std::function done; size_t outstanding = 0; bool settled = false; - /// Keyed by the hashes the answer carried, so agreement is a count rather than a comparison of - /// every answer against every other. - std::unordered_map heard; }; -namespace { - -/// What two members have to agree on before their answer is taken. -/// -/// The hashes, sorted: the storage server returns messages in its own order, and two members -/// holding the same config are not obliged to list it identically. -std::string answer_digest(const retrieved_namespace& got) { - std::vector hashes; - hashes.reserve(got.messages.size()); - for (const auto& m : got.messages) - hashes.emplace_back(m.hash); - std::ranges::sort(hashes); - return "{}"_format(fmt::join(hashes, "\n")); -} - -} // namespace - void Core::_poll() { // Non-owning: the Network is ours alone, and callbacks below must not keep it alive -- doing so // could make the loop thread the last owner and run ~Network there. @@ -299,15 +278,6 @@ void Core::_poll() { }); } -/// How many members have to give the same answer before it is taken without waiting for the rest. -/// -/// Two, which is the smallest number that is more than one member's word -- and it gates *having* -/// the config, not the wait. **Absence never settles this fetch early**, however many members -/// report it: presence is not a majority property. The case this whole thing exists for is a -/// config that has reached one member of a swarm and not the others, so counting "I do not have it" -/// against the one member that does would be the original fault with more steps in it. -static constexpr size_t PROFILE_FETCH_QUORUM = 2; - void Core::fetch_user_profile(std::function done) { // Non-owning: the Network is ours alone, and callbacks below must not keep it alive -- doing so // could make the loop thread the last owner and run ~Network there. @@ -395,7 +365,7 @@ std::vector Core::_profile_retrieve_body() { return to_vector(nlohmann::json{{"requests", std::move(requests)}}.dump()); } -/// One member's answer, tallied against the others. +/// One member's answer. The first that carries a config ends the fetch. void Core::_handle_profile_response( ProfileFanOut& state, const network::service_node& node, @@ -413,18 +383,9 @@ void Core::_handle_profile_response( if (it != json.end() && it->is_array() && !it->empty()) { if (auto got = decode_retrieved( (*it)[0], static_cast(config::Namespace::UserProfile))) { - auto& entry = state.heard[answer_digest(*got)]; - if (entry.agreed == 0) { - entry.node = node.remote_pubkey; - entry.answer = std::move(*got); - } - entry.agreed += 1; - - // Seconded, and carrying something: nothing better is going to arrive, so the - // remaining members are not waited on. An empty answer never takes this exit - // however many members give it -- see PROFILE_FETCH_QUORUM. - if (!entry.answer.messages.empty() && entry.agreed >= PROFILE_FETCH_QUORUM) { - _settle_profile_fetch(state, &entry); + if (!got->messages.empty()) { + ProfileAnswer taken{node.remote_pubkey, std::move(*got)}; + _settle_profile_fetch(state, &taken); return; } } @@ -440,28 +401,13 @@ void Core::_handle_profile_response( timed_out ? "timed out" : "request failed"); } - if (state.outstanding > 0) - return; - - // Everyone has answered and no answer was seconded. **Any member that had the config still - // wins**, because the alternative is not a safer answer but no answer at all -- and the - // ordinary poll, which is what runs instead, believes a single member without asking anyone - // else. Members reporting nothing are not counted against it at all: they are the condition - // being routed around, not evidence. - // - // Among answers that do carry something, the most-agreed wins, and among equals the one - // carrying the most messages -- a member holding part of a config is further along than one - // holding less of it. - ProfileAnswer* best = nullptr; - for (auto& [digest, entry] : state.heard) { - if (entry.answer.messages.empty()) - continue; - if (!best || entry.agreed > best->agreed || - (entry.agreed == best->agreed && - entry.answer.messages.size() > best->answer.messages.size())) - best = &entry; - } - _settle_profile_fetch(state, best); + // Nothing from this one. An empty answer never ends the fetch, however many members give it: + // the case this exists for is a config that has reached one member and not the others, so a + // member saying it has nothing is the condition being routed around rather than evidence + // against the member that has it. Only when every member has answered is there nothing left + // to wait for. + if (state.outstanding == 0) + _settle_profile_fetch(state, nullptr); } /// Ends the fan-out, merging `taken` if there is one. Called exactly once per fetch. From 11b56c6154c98262b74ad85366d51b7500b56dd1 Mon Sep 17 00:00:00 2001 From: Morgan Pretty Date: Wed, 16 Sep 2026 08:24:29 +1000 Subject: [PATCH 3/6] Format the retrieve decoder to the project style `NamespaceIndentation: Inner` indents an anonymous namespace nested inside `session`, and the decoder moved into that namespace without being re-indented. Whitespace and comment rewrapping only: `./utils/format.sh`'s own output, which is what the lint stage runs. --- src/core.cpp | 111 ++++++++++++++++++++++++++------------------------- 1 file changed, 56 insertions(+), 55 deletions(-) diff --git a/src/core.cpp b/src/core.cpp index 7e3cdd9d..a725e59d 100644 --- a/src/core.cpp +++ b/src/core.cpp @@ -173,67 +173,68 @@ struct retrieved_namespace { namespace { -/// Decodes one result of a batch retrieve, or nothing where the node did not answer it. -/// -/// **Nothing and an empty answer are different**, and the difference is load-bearing: a namespace -/// that failed or answered malformedly is not reported to its handler at all, while one that -/// answered with nothing is -- "we asked and there is nothing" is an answer, and some handlers act -/// on it. -/// -/// Takes no node and touches no database, deliberately. Which node an answer came from matters to -/// whoever called: the retrieve cursor is kept per node and has to stay with the one that produced -/// it, and a decoder that knew about nodes is the shape in which one node's cursor gets written -/// from another node's response -- silently, because each node's cursor is individually plausible. -std::optional decode_retrieved(const nlohmann::json& res, int16_t ns_val) { - auto code_it = res.find("code"); - if (code_it == res.end() || code_it->get() != 200) { - log::warning(cat, "Retrieve of namespace {} failed: {}", ns_val, res.dump()); - return std::nullopt; - } - auto body_it = res.find("body"); - if (body_it == res.end()) - return std::nullopt; - auto msgs_it = body_it->find("messages"); - if (msgs_it == body_it->end() || !msgs_it->is_array()) - return std::nullopt; - - retrieved_namespace got; - if (auto m = body_it->find("more"); m != body_it->end() && m->is_boolean()) - got.more = m->get(); - - log::debug(cat, "Retrieved {} message(s) from namespace {}", msgs_it->size(), ns_val); - - for (const auto& msg : *msgs_it) { - auto data_it = msg.find("data"); - if (data_it == msg.end() || !data_it->is_string()) - continue; - auto& decoded = got.data.emplace_back(); - auto b64 = data_it->get(); - decoded.reserve(oxenc::from_base64_size(b64.size())); - oxenc::from_base64(b64.begin(), b64.end(), std::back_inserter(decoded)); + /// Decodes one result of a batch retrieve, or nothing where the node did not answer it. + /// + /// **Nothing and an empty answer are different**, and the difference is load-bearing: a + /// namespace that failed or answered malformedly is not reported to its handler at all, while + /// one that answered with nothing is -- "we asked and there is nothing" is an answer, and some + /// handlers act on it. + /// + /// Takes no node and touches no database, deliberately. Which node an answer came from matters + /// to whoever called: the retrieve cursor is kept per node and has to stay with the one that + /// produced it, and a decoder that knew about nodes is the shape in which one node's cursor + /// gets written from another node's response -- silently, because each node's cursor is + /// individually plausible. + std::optional decode_retrieved(const nlohmann::json& res, int16_t ns_val) { + auto code_it = res.find("code"); + if (code_it == res.end() || code_it->get() != 200) { + log::warning(cat, "Retrieve of namespace {} failed: {}", ns_val, res.dump()); + return std::nullopt; + } + auto body_it = res.find("body"); + if (body_it == res.end()) + return std::nullopt; + auto msgs_it = body_it->find("messages"); + if (msgs_it == body_it->end() || !msgs_it->is_array()) + return std::nullopt; + + retrieved_namespace got; + if (auto m = body_it->find("more"); m != body_it->end() && m->is_boolean()) + got.more = m->get(); + + log::debug(cat, "Retrieved {} message(s) from namespace {}", msgs_it->size(), ns_val); + + for (const auto& msg : *msgs_it) { + auto data_it = msg.find("data"); + if (data_it == msg.end() || !data_it->is_string()) + continue; + auto& decoded = got.data.emplace_back(); + auto b64 = data_it->get(); + decoded.reserve(oxenc::from_base64_size(b64.size())); + oxenc::from_base64(b64.begin(), b64.end(), std::back_inserter(decoded)); - SwarmMessage swarm_msg; - swarm_msg.data = {decoded.data(), decoded.size()}; + SwarmMessage swarm_msg; + swarm_msg.data = {decoded.data(), decoded.size()}; - if (auto h = msg.find("hash"); h != msg.end() && h->is_string()) - swarm_msg.hash = h->get(); + if (auto h = msg.find("hash"); h != msg.end() && h->is_string()) + swarm_msg.hash = h->get(); - if (auto t = msg.find("timestamp"); t != msg.end() && t->is_number_integer()) - swarm_msg.timestamp = from_epoch_ms(t->get()); + if (auto t = msg.find("timestamp"); t != msg.end() && t->is_number_integer()) + swarm_msg.timestamp = from_epoch_ms(t->get()); - if (auto e = msg.find("expiry"); e != msg.end() && e->is_number_integer()) - swarm_msg.expiry = from_epoch_ms(e->get()); + if (auto e = msg.find("expiry"); e != msg.end() && e->is_number_integer()) + swarm_msg.expiry = from_epoch_ms(e->get()); - got.messages.push_back(std::move(swarm_msg)); - } + got.messages.push_back(std::move(swarm_msg)); + } - // A node claiming more while returning nothing cannot be continued: there is no new hash to - // move the cursor to, so another round would ask the same question and get the same answer. - // Reported as finished instead, or a handler waiting on `is_final` would wait for one that - // never comes. - got.more = got.more && !got.messages.empty(); - return got; -} + // A node claiming more while returning nothing cannot be continued: there is no new hash to + // move the cursor to, so another round would ask the same question and get the same answer. + // Reported as finished instead, or a handler waiting on `is_final` would wait for one that + // never comes. + got.more = got.more && !got.messages.empty(); + return got; + } } // namespace From a6565cdd89ffae7bfcaf544248a1b0cba3992a4c Mon Sep 17 00:00:00 2001 From: Morgan Pretty Date: Wed, 23 Sep 2026 10:02:48 +1000 Subject: [PATCH 4/6] Queue the profile fan-out on Core's job queue The swarm lookup and every member's answer were marshalled onto `_loop`. Core now has its own `_jq`, and work queued there is cancelled when Core goes away, so an answer that lands during teardown is dropped rather than run against components that are already being destroyed. That matters more with the next change, where late answers do real work. --- include/session/core.hpp | 4 ++-- src/core.cpp | 23 +++++++++++++---------- 2 files changed, 15 insertions(+), 12 deletions(-) diff --git a/include/session/core.hpp b/include/session/core.hpp index 891dbb92..760050ec 100644 --- a/include/session/core.hpp +++ b/include/session/core.hpp @@ -581,8 +581,8 @@ class Core { /// protection against something that already repairs itself, and costs a round trip in front of /// somebody watching a progress indicator. /// - /// `done` is called exactly once, on Core's loop, with whether a config was merged. Safe to - /// call without waiting on it: what it finds is merged into `configs` like anything a poll + /// `done` is called exactly once, on Core's job queue, with whether a config was merged. Safe + /// to call without waiting on it: what it finds is merged into `configs` like anything a poll /// brings in, and the retrieve cursor it records means the next poll carries on from there. void fetch_user_profile(std::function done); diff --git a/src/core.cpp b/src/core.cpp index a725e59d..65ce911f 100644 --- a/src/core.cpp +++ b/src/core.cpp @@ -249,8 +249,8 @@ struct Core::ProfileAnswer { /// How much of a fan-out is still outstanding. /// -/// Shared between every in-flight request and touched only on Core's loop, which is what makes a -/// plain count safe here: the responses arrive on the network's threads and are marshalled across +/// Shared between every in-flight request and touched only on Core's job queue, which is what makes +/// a plain count safe here: the responses arrive on the network's threads and are marshalled across /// before any of this is read. struct Core::ProfileFanOut { std::function done; @@ -293,7 +293,10 @@ void Core::fetch_user_profile(std::function done) { state->done = std::move(done); net->get_swarm(globals.pubkey_x25519(), false, [this, net, state](auto, auto swarm) { - _loop.call([this, net, state, swarm = std::move(swarm)] { + // Onto our own queue rather than the loop's: work queued here is cancelled when Core goes + // away, so a swarm lookup or an answer landing during teardown is dropped instead of run + // against components that are already being destroyed. + _jq.call([this, net, state, swarm = std::move(swarm)] { if (state->settled) return; if (swarm.empty()) { @@ -319,13 +322,13 @@ void Core::fetch_user_profile(std::function done) { std::optional body) { // Marshalled rather than handled here: these arrive on the network's // threads, several at once by design, and everything they touch -- - // the tally, the config merge, the cursor -- belongs to Core's loop. - _loop.call([this, - state, - node, - success, - timeout, - body = std::move(body)]() mutable { + // the tally, the config merge, the cursor -- belongs to Core's queue. + _jq.call([this, + state, + node, + success, + timeout, + body = std::move(body)]() mutable { _handle_profile_response( *state, node, From 8716bd4d7c28b8276a341a8e8fbd9f21f900815e Mon Sep 17 00:00:00 2001 From: Morgan Pretty Date: Wed, 23 Sep 2026 10:03:02 +1000 Subject: [PATCH 5/6] Merge every member's config, and report after a short window The fetch settled on the first member to answer with a config: it merged that one, called `done`, and dropped every answer after it. Two problems with that. A member can be behind its swarm, so the first config to arrive is not necessarily the newest -- and a restore that lands on a stale member reports a stale profile to the user watching it. And the answers dropped afterwards were real data from our own swarm, which the next ordinary poll then fetched again. So absorbing an answer and concluding the fetch are now separate. Every answer that carries a config is merged and has its member's cursor recorded, whenever it arrives, including after `done` has been called. The first such answer opens a 500ms window instead of ending the fetch, and the fetch concludes when every member has answered or the window closes, whichever is first -- `done` exactly once. Configs merge rather than replace, so each extra answer can only bring the result forward. An empty answer still never ends the fetch on its own, and only a config opens the window: an answer with nothing in it is the condition the fan-out exists to route around. Concluding no longer takes an answer, so the pointer passed to it -- and the struct that carried a decoded answer to it -- are gone. The decode moves out of the merge's `try`, so a failure while merging is not logged as a failure to parse. Four tests, where there were none: a stale first config overtaken inside the window, a late config merged after the fetch reported with its cursor written against the member that gave it, the window closing early once everyone has answered, and all-empty concluding with nothing found. Run against the previous behaviour, the first three fail in seven places; the fourth passes on both, since that behaviour is unchanged. `MockNetwork` can now hand back a whole swarm. --- include/session/core.hpp | 32 +++--- src/core.cpp | 123 ++++++++++---------- tests/CMakeLists.txt | 1 + tests/test_helper.hpp | 9 +- tests/test_profile_fetch.cpp | 214 +++++++++++++++++++++++++++++++++++ 5 files changed, 307 insertions(+), 72 deletions(-) create mode 100644 tests/test_profile_fetch.cpp diff --git a/include/session/core.hpp b/include/session/core.hpp index 760050ec..f68faf17 100644 --- a/include/session/core.hpp +++ b/include/session/core.hpp @@ -358,17 +358,18 @@ class Core { std::span messages); // The fan-out behind fetch_user_profile(). Defined in core.cpp: nothing outside it needs the - // shape, and one of them holds a decoded response. - struct ProfileAnswer; + // shape. struct ProfileFanOut; std::vector _profile_retrieve_body(); void _handle_profile_response( - ProfileFanOut& state, + const std::shared_ptr& state, const network::service_node& node, std::optional body, bool timed_out); - void _settle_profile_fetch(ProfileFanOut& state, ProfileAnswer* taken); + void _absorb_profile_answer( + const network::ed25519_pubkey& node, std::span messages); + void _conclude_profile_fetch(ProfileFanOut& state); // Decrypts and dispatches one-to-one messages from Namespace::Default. void _handle_direct_messages(std::span messages); @@ -571,15 +572,20 @@ class Core { /// Asks every member concurrently, and is deliberately asymmetric about what it will conclude /// from what they say. **Having the config is not a majority property**: one member holding it /// is the whole answer, and the others not having it yet is the condition being routed around - /// rather than evidence against it. So the first answer carrying a config ends the fetch, and - /// an empty answer never does however many members give it -- only every member having answered - /// concludes that there is nothing to find. - /// - /// **No agreement is required of the answer that wins**, on purpose. A config is merged rather - /// than assigned, so a stale one taken from a member behind its swarm is corrected by the next - /// ordinary poll rather than stuck; waiting for a second member to say the same thing buys - /// protection against something that already repairs itself, and costs a round trip in front of - /// somebody watching a progress indicator. + /// rather than evidence against it. So an empty answer never ends the fetch however many + /// members give it -- only every member having answered concludes that there is nothing to + /// find. + /// + /// **The first config opens a short window rather than ending the fetch.** A member can be + /// behind its swarm, so the first config to arrive may not be the newest; whatever else arrives + /// over the next half second is merged before `done` is called, and since configs merge rather + /// than replace, each extra answer can only bring the result forward. No agreement is sought + /// -- nothing waits for a second member to say the same thing -- and the window is short + /// because somebody is watching a progress indicator. It closes early once every member has + /// answered. + /// + /// **Answers keep being merged after `done`.** One arriving late still came from our own + /// swarm, and dropping it would only leave the next ordinary poll to fetch it again. /// /// `done` is called exactly once, on Core's job queue, with whether a config was merged. Safe /// to call without waiting on it: what it finds is merged into `configs` like anything a poll diff --git a/src/core.cpp b/src/core.cpp index 65ce911f..032bd1b0 100644 --- a/src/core.cpp +++ b/src/core.cpp @@ -238,24 +238,27 @@ namespace { } // namespace -/// The answer a fetch settled on, and the member that gave it. +/// How long the first member's config is left open for others to add to. /// -/// The member is kept with it because the retrieve cursor is written against the node that -/// produced the messages and no other. -struct Core::ProfileAnswer { - network::ed25519_pubkey node; - retrieved_namespace answer; -}; - -/// How much of a fan-out is still outstanding. +/// A member can be behind its swarm, so the first config to arrive is not necessarily the newest. +/// Merging whatever else comes in over a short window lets a stale first answer be overtaken before +/// the fetch reports -- and configs merge rather than replace, so each extra answer can only bring +/// the result forward. Short because somebody is watching a progress indicator: long enough to let +/// in a few more members, not to wait for all of them. +constexpr std::chrono::milliseconds PROFILE_FETCH_GRACE{500}; + +/// How far a fan-out has got. /// -/// Shared between every in-flight request and touched only on Core's job queue, which is what makes -/// a plain count safe here: the responses arrive on the network's threads and are marshalled across -/// before any of this is read. +/// Shared between every in-flight request and the grace timer, and touched only on Core's job +/// queue, which is what makes plain fields safe here: the responses arrive on the network's threads +/// and are marshalled across before any of this is read. struct Core::ProfileFanOut { std::function done; size_t outstanding = 0; - bool settled = false; + /// Whether any member has answered with a config, which is what `done` reports. + bool found = false; + bool grace_started = false; + bool concluded = false; }; void Core::_poll() { @@ -297,14 +300,14 @@ void Core::fetch_user_profile(std::function done) { // away, so a swarm lookup or an answer landing during teardown is dropped instead of run // against components that are already being destroyed. _jq.call([this, net, state, swarm = std::move(swarm)] { - if (state->settled) + if (state->concluded) return; if (swarm.empty()) { // Not an error worth failing loudly over: the ordinary poll still runs, and this // was only ever an attempt to get there sooner. A swarm of one is still worth // asking -- one member holding the config is the whole answer. log::warning(cat, "Cannot fan out a profile fetch: no swarm members available"); - _settle_profile_fetch(*state, nullptr); + _conclude_profile_fetch(*state); return; } @@ -330,7 +333,7 @@ void Core::fetch_user_profile(std::function done) { timeout, body = std::move(body)]() mutable { _handle_profile_response( - *state, + state, node, (success && body) ? std::move(body) : std::nullopt, timeout); @@ -369,31 +372,28 @@ std::vector Core::_profile_retrieve_body() { return to_vector(nlohmann::json{{"requests", std::move(requests)}}.dump()); } -/// One member's answer. The first that carries a config ends the fetch. +/// One member's answer, whenever it arrives -- including after the fetch has reported. +/// +/// **Every config that arrives is merged**, even one landing after `done` has been called: it came +/// from a member of our own swarm, and dropping it would only leave the next ordinary poll to fetch +/// it again. What the end of the fetch decides is when `done` is called, not which answers count. void Core::_handle_profile_response( - ProfileFanOut& state, + const std::shared_ptr& state, const network::service_node& node, std::optional body, bool timed_out) { - if (state.settled) - return; - - state.outstanding -= 1; + --state->outstanding; + // Decoded here and merged below, outside the `try`: a failure while merging is not a failure to + // parse, and reporting it as one would send somebody looking at the wrong thing. + std::optional got; if (body) { try { auto json = nlohmann::json::parse(*body); auto it = json.find("results"); - if (it != json.end() && it->is_array() && !it->empty()) { - if (auto got = decode_retrieved( - (*it)[0], static_cast(config::Namespace::UserProfile))) { - if (!got->messages.empty()) { - ProfileAnswer taken{node.remote_pubkey, std::move(*got)}; - _settle_profile_fetch(state, &taken); - return; - } - } - } + if (it != json.end() && it->is_array() && !it->empty()) + got = decode_retrieved( + (*it)[0], static_cast(config::Namespace::UserProfile)); } catch (const std::exception& e) { log::warning(cat, "Failed to parse profile fetch response: {}", e.what()); } @@ -405,40 +405,47 @@ void Core::_handle_profile_response( timed_out ? "timed out" : "request failed"); } - // Nothing from this one. An empty answer never ends the fetch, however many members give it: - // the case this exists for is a config that has reached one member and not the others, so a - // member saying it has nothing is the condition being routed around rather than evidence - // against the member that has it. Only when every member has answered is there nothing left - // to wait for. - if (state.outstanding == 0) - _settle_profile_fetch(state, nullptr); -} + if (got && !got->messages.empty()) { + _absorb_profile_answer(node.remote_pubkey, got->messages); + state->found = true; + // The first config opens the window rather than ending the fetch; see PROFILE_FETCH_GRACE. + // Only a config opens it: an answer with nothing in it is the condition this routes around. + if (!state->concluded && !state->grace_started) { + state->grace_started = true; + _jq.call_later(PROFILE_FETCH_GRACE, [this, state] { _conclude_profile_fetch(*state); }); + } + } -/// Ends the fan-out, merging `taken` if there is one. Called exactly once per fetch. -void Core::_settle_profile_fetch(ProfileFanOut& state, ProfileAnswer* taken) { - if (state.settled) - return; - state.settled = true; + // An empty answer never ends the fetch on its own; the last one to arrive does. + if (state->outstanding == 0) + _conclude_profile_fetch(*state); +} - bool found = taken && !taken->answer.messages.empty(); +/// Merges one member's config and records where that member's next retrieve resumes. +void Core::_absorb_profile_answer( + const network::ed25519_pubkey& node, std::span messages) { + auto configs_held = configs.batch(); + // Final, even where the member said it was holding more. This is a one-shot fetch rather than + // a poll that continues: the cursor recorded below is what lets the ordinary poll pick up + // anything left, and a handler told to wait for a final that never comes would wait for ever. + receive_messages(messages, config::Namespace::UserProfile, true); + _record_swarm_cursor(node, config::Namespace::UserProfile, messages); +} - if (taken && !taken->answer.messages.empty()) { - auto configs_held = configs.batch(); - // Final, even where the member said it was holding more. This is a one-shot fetch rather - // than a poll that continues: the cursor recorded below is what lets the ordinary poll pick - // up anything left, and a handler told to wait for a final that never comes would wait for - // ever. - receive_messages(taken->answer.messages, config::Namespace::UserProfile, true); - _record_swarm_cursor(taken->node, config::Namespace::UserProfile, taken->answer.messages); - } +/// Reports the fan-out's result, once. Whichever of the grace timer and the last answer comes +/// first reports it; the other finds it already done. +void Core::_conclude_profile_fetch(ProfileFanOut& state) { + if (state.concluded) + return; + state.concluded = true; log::info( cat, - "Profile fetch settled: {}", - found ? "config merged" : "nothing held by the swarm"); + "Profile fetch concluded: {}", + state.found ? "config merged" : "nothing held by the swarm"); if (state.done) - state.done(found); + state.done(state.found); } void Core::_send_poll( diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index ed631566..7a2c557e 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -69,6 +69,7 @@ list(APPEND LIB_SESSION_UTESTS_SOURCES test_snode_pool.cpp test_core_network.cpp test_poll.cpp + test_profile_fetch.cpp test_pfs_key_cache.cpp test_sqlite_bind.cpp) diff --git a/tests/test_helper.hpp b/tests/test_helper.hpp index 08e5c2dd..47122a6f 100644 --- a/tests/test_helper.hpp +++ b/tests/test_helper.hpp @@ -46,6 +46,10 @@ class MockNetwork : public network::Network { // The node returned by get_swarm; tests can change this to simulate swarm-member switches. network::service_node current_node; + // The whole swarm, for code that asks every member rather than one. Returned instead of + // `current_node` when set. + std::vector current_swarm; + void send_request( network::Request request, network::network_response_callback_t callback) override { sent_requests.push_back({std::move(request), std::move(callback)}); @@ -57,7 +61,10 @@ class MockNetwork : public network::Network { std::function< void(network::swarm_id_t swarm_id, std::vector swarm)> callback) override { - callback(0, {current_node}); + if (current_swarm.empty()) + callback(0, {current_node}); + else + callback(0, current_swarm); } std::vector downloads; diff --git a/tests/test_profile_fetch.cpp b/tests/test_profile_fetch.cpp new file mode 100644 index 00000000..cabb4cbe --- /dev/null +++ b/tests/test_profile_fetch.cpp @@ -0,0 +1,214 @@ +#include + +#include +#include +#include +#include +#include +#include + +#include "test_helper.hpp" + +using namespace session; +using namespace std::literals; + +namespace { + +/// Longer than the fetch's grace window, so a wait of this long has seen it close. +constexpr auto PAST_THE_WINDOW = 800ms; + +constexpr auto USER_PROFILE = static_cast(config::Namespace::UserProfile); + +/// A swarm of `n` members with distinct keys, so each keeps a retrieve cursor of its own. +std::vector swarm_of(size_t n) { + std::vector swarm(n); + for (size_t i = 0; i < n; i++) + swarm[i].remote_pubkey[0] = std::byte(i + 1); + return swarm; +} + +/// Another device of the same account, descending from what this one has already published: a +/// device built from nothing would land on our own seqno and be a merge of two histories, which is +/// not what this fetch is about. +/// +/// Our side is read on Core's loop, the only thread allowed to touch its configs. The other device +/// is a plain config object of the test's own, and needs no such care. +config::UserProfile another_device(TempCore& c) { + auto dump = TestHelper::on_loop(*c, [&] { + auto& ours = c->configs.user_profile(); + auto [seqno, messages, obsolete] = ours.push(); + ours.confirm_pushed(seqno, {"seededprofile"}); + c->configs.store_dumps(); + return ours.make_dump(); + }); + auto seed = c->globals.account_seed(); + return config::UserProfile{seed.ed25519_secret(), dump}; +} + +/// The profile name this Core holds, read where it is allowed to be read and copied out. +std::optional name_of(TempCore& c) { + return TestHelper::on_loop(*c, [&]() -> std::optional { + auto name = c->configs.user_profile().get_name(); + return name ? std::optional{std::string{*name}} : std::nullopt; + }); +} + +/// What that device pushes after renaming itself. +std::vector> renamed(config::UserProfile& theirs, std::string_view name) { + theirs.set_name(name); + auto [seqno, messages, obsolete] = theirs.push(); + theirs.confirm_pushed(seqno, {"pushed{}"_format(seqno)}); + return messages; +} + +/// A member's answer to the fetch's single retrieve, carrying `messages` under `hash`. +std::string answer(const std::vector>& messages, std::string_view hash) { + nlohmann::json body; + body["messages"] = nlohmann::json::array(); + for (const auto& m : messages) + body["messages"].push_back({{"data", oxenc::to_base64(m)}, {"hash", hash}}); + return nlohmann::json{{"results", {{{"code", 200}, {"body", std::move(body)}}}}}.dump(); +} + +std::string empty_answer() { + return answer({}, ""); +} + +/// The fetch's request to each member, in the order they were sent. +std::vector sent_to_each(MockNetwork& net) { + return std::exchange(net.sent_requests, {}); +} + +void reply(TempCore& c, MockNetwork::SentRequest& to, std::string body) { + to.callback(true, false, 200, {}, std::move(body)); + TestHelper::drain(*c); +} + +} // namespace + +TEST_CASE("Profile fetch: a stale first config is overtaken inside the window", "[core][profile]") { + TempCore c; + auto* net = attach_mock_network(*c); + net->current_swarm = swarm_of(3); + auto theirs = another_device(c); + auto older = renamed(theirs, "Stale"); + auto newer = renamed(theirs, "Current"); + + int calls = 0; + bool found = false; + c->fetch_user_profile([&](bool f) { + calls++; + found = f; + }); + TestHelper::drain(*c); + auto sent = sent_to_each(*net); + REQUIRE(sent.size() == 3); + + // The first member is behind its swarm. Its config is merged, and the fetch does not report on + // it: that is the whole point of the window. + reply(c, sent[0], answer(older, "h-old")); + CHECK(name_of(c) == "Stale"); + CHECK(calls == 0); + + // A second member has the newer one, inside the window. + reply(c, sent[1], answer(newer, "h-new")); + CHECK(name_of(c) == "Current"); + CHECK(calls == 0); + + // The third never answers; the window closes on its own. + std::this_thread::sleep_for(PAST_THE_WINDOW); + TestHelper::drain(*c); + CHECK(calls == 1); + CHECK(found); + CHECK(name_of(c) == "Current"); +} + +TEST_CASE("Profile fetch: a config arriving after it reported is still merged", "[core][profile]") { + TempCore c; + auto* net = attach_mock_network(*c); + net->current_swarm = swarm_of(2); + auto theirs = another_device(c); + auto older = renamed(theirs, "Stale"); + auto newer = renamed(theirs, "Current"); + + int calls = 0; + c->fetch_user_profile([&](bool) { calls++; }); + TestHelper::drain(*c); + auto sent = sent_to_each(*net); + REQUIRE(sent.size() == 2); + + reply(c, sent[0], answer(older, "h-old")); + std::this_thread::sleep_for(PAST_THE_WINDOW); + TestHelper::drain(*c); + REQUIRE(calls == 1); + + // Too late to be reported, not too late to count: it came from our own swarm, and dropping it + // would only leave the next poll to fetch it again. + reply(c, sent[1], answer(newer, "h-new")); + CHECK(name_of(c) == "Current"); + CHECK(calls == 1); + + // And the cursor is written against the member that gave it, not the one that was first. + CHECK(TestHelper::namespace_last_hash(*c, USER_PROFILE, net->current_swarm[1].remote_pubkey) == + "h-new"); + CHECK(TestHelper::namespace_last_hash(*c, USER_PROFILE, net->current_swarm[0].remote_pubkey) == + "h-old"); +} + +TEST_CASE( + "Profile fetch: the window closes early once every member has answered", + "[core][profile]") { + TempCore c; + auto* net = attach_mock_network(*c); + net->current_swarm = swarm_of(2); + auto theirs = another_device(c); + auto config = renamed(theirs, "Current"); + + int calls = 0; + bool found = false; + c->fetch_user_profile([&](bool f) { + calls++; + found = f; + }); + TestHelper::drain(*c); + auto sent = sent_to_each(*net); + REQUIRE(sent.size() == 2); + + reply(c, sent[0], answer(config, "h1")); + CHECK(calls == 0); + // Nobody is left to wait for, so there is no reason to sit out the rest of the window. + reply(c, sent[1], empty_answer()); + CHECK(calls == 1); + CHECK(found); + + // Nor does the window's own end report a second time. + std::this_thread::sleep_for(PAST_THE_WINDOW); + TestHelper::drain(*c); + CHECK(calls == 1); +} + +TEST_CASE("Profile fetch: every member empty concludes with nothing found", "[core][profile]") { + TempCore c; + auto* net = attach_mock_network(*c); + net->current_swarm = swarm_of(3); + + int calls = 0; + bool found = true; + c->fetch_user_profile([&](bool f) { + calls++; + found = f; + }); + TestHelper::drain(*c); + auto sent = sent_to_each(*net); + REQUIRE(sent.size() == 3); + + // An empty answer never ends the fetch on its own, however many give it... + reply(c, sent[0], empty_answer()); + reply(c, sent[1], empty_answer()); + CHECK(calls == 0); + + // ...and the last one does, at once: there is no window, since nothing opened one. + reply(c, sent[2], empty_answer()); + CHECK(calls == 1); + CHECK_FALSE(found); +} From 0fbf4b7a439dce87af1cdeb1d301b961ab631f8f Mon Sep 17 00:00:00 2001 From: Morgan Pretty Date: Wed, 23 Sep 2026 10:04:20 +1000 Subject: [PATCH 6/6] Sign the profile retrieve without asking whether to The body builder checked `retrieve_requires_auth` before signing. For the profile namespace that is always true -- it is owner-writable -- and signing is allowed on any namespace, so the branch had no case in which skipping the signature would have been right. It is now an assertion of the invariant, and the signing is unconditional. A test holds the request's shape: every member is asked for the profile namespace, signed, and without a cursor. --- src/core.cpp | 16 +++++++++------- tests/test_profile_fetch.cpp | 24 ++++++++++++++++++++++++ 2 files changed, 33 insertions(+), 7 deletions(-) diff --git a/src/core.cpp b/src/core.cpp index 032bd1b0..335a85a1 100644 --- a/src/core.cpp +++ b/src/core.cpp @@ -8,6 +8,7 @@ #include #include +#include #include #include #include @@ -359,13 +360,14 @@ std::vector Core::_profile_retrieve_body() { {"namespace", ns_val}, }; - if (retrieve_requires_auth(ns_val)) { - auto seed = globals.account_seed(); - auto to_sign = ns_signature_value("retrieve", ns_val, now_ms); - params["pubkey_ed25519"] = globals.pubkey_ed25519().hex(); - params["timestamp"] = now_ms; - params["signature"] = "{:b}"_format(ed25519::sign(seed.ed25519_secret(), to_span(to_sign))); - } + // The profile is owner-writable, so retrieving it always needs a signature -- and signing is + // allowed on any namespace, so there is no case in which leaving it out would be right. + assert(retrieve_requires_auth(ns_val)); + auto seed = globals.account_seed(); + auto to_sign = ns_signature_value("retrieve", ns_val, now_ms); + params["pubkey_ed25519"] = globals.pubkey_ed25519().hex(); + params["timestamp"] = now_ms; + params["signature"] = "{:b}"_format(ed25519::sign(seed.ed25519_secret(), to_span(to_sign))); nlohmann::json requests = nlohmann::json::array(); requests.push_back({{"method", "retrieve"}, {"params", std::move(params)}}); diff --git a/tests/test_profile_fetch.cpp b/tests/test_profile_fetch.cpp index cabb4cbe..834d0e54 100644 --- a/tests/test_profile_fetch.cpp +++ b/tests/test_profile_fetch.cpp @@ -212,3 +212,27 @@ TEST_CASE("Profile fetch: every member empty concludes with nothing found", "[co CHECK(calls == 1); CHECK_FALSE(found); } + +TEST_CASE("Profile fetch: every member is asked with a signed retrieve", "[core][profile]") { + TempCore c; + auto* net = attach_mock_network(*c); + net->current_swarm = swarm_of(2); + + c->fetch_user_profile([](bool) {}); + TestHelper::drain(*c); + auto sent = sent_to_each(*net); + REQUIRE(sent.size() == 2); + + for (auto& request : sent) { + auto batch = parse_json(*request.request.body); + REQUIRE(batch["requests"].size() == 1); + auto params = batch["requests"][0]["params"]; + CHECK(params["namespace"] == USER_PROFILE); + CHECK(params.contains("signature")); + CHECK(params.contains("pubkey_ed25519")); + CHECK(params.contains("timestamp")); + // No cursor: this asks members that may never have been asked before, so resuming from + // any one member's position would ask the wrong question. + CHECK_FALSE(params.contains("last_hash")); + } +}