From c1d23b5888feb18580eefbf11253b6048abd1134 Mon Sep 17 00:00:00 2001 From: Joe Rivera Date: Tue, 29 Sep 2026 22:35:01 -0500 Subject: [PATCH 1/2] feat(state): collapse HTTP cache entry to a single data() blob, drop cache_mutex (Phase 2) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Phase 2 / atomicity of docs/roadmap/STATE_LIFETIME_AND_ATOMICITY.md, building on the Phase 1 owning accessors (#492). Before: an entry was 7 metadata child value-nodes (status/etag/last_modified/ content_type/effective_url/fetched_at/freshness_ttl) plus the body on the entry node's data(). Reading/writing that spread across N node locks needed a process-wide cache_mutex to avoid a torn (half-written) entry AND to stop a concurrent sweep freeing a node under a bare pointer. After: the whole entry is ONE opaque blob on the entry node's data() (version byte + 8 length-prefixed fields). A read or write is now a SINGLE atomic node op under that node's own _mutex — a reader sees the whole old blob or the whole new blob, never a mix. read_entry uses findDescendantShared() (Phase 1) to PIN the node, so a concurrent sweepExpired() can only unlink it, never free it under the reader. Writers of one key are already serialized by the single-flight flight_mutex, and the write is atomic, so there is nothing left for a global lock to protect: cache_mutex is REMOVED. Reads no longer take any process-wide lock (lock-free apart from the node's own brief _mutex). Net: lock-free, torn-read-free cache reads; the cache's only write coordination is the pre-existing per-key single-flight. Entry child nodes are no longer individually addressable (they were internal-only; no consumer read them — verified). Test file unchanged (its checks are entry presence/absence + expireAt, both preserved). All 12 HttpCache tests pass, incl. ConcurrentDistinctKeysAreServedSafely now without cache_mutex; AriadneStateUri/StateLifetime/AriadneNetIntrinsics green (38/38). --- src/cvc/ariadne/uri_http_cache.cpp | 174 +++++++++++++++++------------ 1 file changed, 101 insertions(+), 73 deletions(-) diff --git a/src/cvc/ariadne/uri_http_cache.cpp b/src/cvc/ariadne/uri_http_cache.cpp index d7201f3a..82b67703 100644 --- a/src/cvc/ariadne/uri_http_cache.cpp +++ b/src/cvc/ariadne/uri_http_cache.cpp @@ -13,15 +13,15 @@ // evicts and frees the body blob — one mechanism, the state tree's own expiry. // - single-flight = concurrent requests for one URL coalesce to a single transfer. // -// CONCURRENCY: cvc::state uses per-node locks and findDescendant()/operator() hand back a bare -// pointer/reference whose only owner is the parent's child map, so a concurrent sweepExpired() can -// free a node mid-access. Every touch of the `entries` subtree (sweep + snapshot read + write + -// invalidation) is therefore serialized under one process-wide `cache_mutex`, held ONLY for the -// brief tree work and released across the (blocking) network send — so reads never overlap eviction -// and a multi-node entry is written/read as a unit. (Assumes nothing subscribes to the cache -// subtree and re-enters a resolve from sweepExpired()'s `expiring`/`childChanged` signals — the -// cache nodes are internal.) `cache_mutex` and the single-flight `flight_mutex` are never held -// simultaneously. +// CONCURRENCY: each entry is stored as ONE opaque blob on its node's data() (never as separate +// metadata child nodes), so a read or write is a SINGLE atomic node operation under that node's own +// _mutex — a reader sees the whole old blob or the whole new blob, never a torn mix. Lifetime is +// covered by findDescendantShared(): read_entry PINS the entry node, so a concurrent sweepExpired() +// can only UNLINK it, never free it under the reader (cvc::state's owning-accessor lifetime fix). +// Writers of one key are already serialized by the single-flight `flight_mutex`, and the retention +// sweep only unlinks pinned-safe nodes. Together these remove the need for any process-wide cache +// lock: there is no `cache_mutex`. (The cache nodes are internal, so nothing re-enters a resolve +// from sweepExpired()'s `expiring`/`childChanged` signals.) // // Clock is UTC wall time (boost::posix_time::microsec_clock::universal_time()) — the same clock the // node-expiry path uses, and correct for HTTP freshness (never the sim-time world_clock). A @@ -35,6 +35,7 @@ #include #include #include +#include #include #include #include @@ -62,12 +63,6 @@ constexpr long kDefaultFreshnessSecs = 60; // fresh window when no Cache-Control constexpr long kRetentionSecs = 3600; // hard-eviction floor (retention >= freshness) constexpr char kCacheRoot[] = "sys.net.http_cache.entries"; -// One process-wide lock serializing all access to the cache subtree (see the CONCURRENCY note). -std::mutex &cache_mutex() { - static std::mutex m; - return m; -} - // --- URL normalization + key ------------------------------------------------------------------- // True if the URL authority carries a `user[:pass]@` userinfo component. Such a URL is credential- @@ -244,10 +239,9 @@ long derive_freshness_ttl(const std::vector &headers, const CacheCo return has_validator ? 0 : kDefaultFreshnessSecs; } -// --- entry read/write (all callers hold cache_mutex) ------------------------------------------- +// --- entry read/write (single-node blob; no cache lock) ---------------------------------------- -// A snapshot of a cached entry, captured as a value copy so nothing downstream holds a bare node -// pointer past the lock. +// A decoded snapshot of a cached entry (a value copy — nothing downstream holds a node pointer). struct EntrySnapshot { bool present = false; std::string body; @@ -259,35 +253,83 @@ struct EntrySnapshot { std::string freshness_ttl; }; -std::string child_value(cvc::state *e, const char *name) { - cvc::state *c = e->findDescendant(name); - return c ? c->value() : std::string(); +// The entry is one opaque blob on the node's data(): a version byte, then 8 length-prefixed fields +// (little-endian uint64 length + bytes) in a fixed order — status, etag, last_modified, +// content_type, effective_url, fetched_at, freshness_ttl, body. Encoding everything into ONE value +// makes a read/write a single atomic node op (no torn multi-node entry), which is what lets the +// cache drop its process-wide lock (see the CONCURRENCY note). +constexpr unsigned char kEntryBlobVersion = 1; + +void put_field(std::string &out, const std::string &f) { + const std::uint64_t n = f.size(); + for (int i = 0; i < 8; ++i) + out.push_back(static_cast((n >> (8 * i)) & 0xFF)); + out.append(f); } -// Read the entry at `key`; present == a node WITH a committed body blob on data(). A metadata-only -// node (a write still in flight elsewhere — though cache_mutex serializes writes, so this is belt- -// and-braces) reads as absent. CALLER MUST HOLD cache_mutex (the bare node pointer is used only -// within the critical section). -EntrySnapshot read_entry(cvc::state &entries, const std::string &key) { +bool get_field(const std::string &in, std::size_t &pos, std::string &out) { + if (pos + 8 > in.size()) + return false; + std::uint64_t n = 0; + for (int i = 0; i < 8; ++i) + n |= static_cast(static_cast(in[pos + i])) << (8 * i); + pos += 8; + if (n > in.size() - pos) + return false; + out.assign(in, pos, static_cast(n)); + pos += static_cast(n); + return true; +} + +std::string encode_entry(long status, const std::string &etag, const std::string &last_modified, + const std::string &content_type, const std::string &effective_url, + const std::string &fetched_at_iso, long freshness_ttl, + const std::string &body) { + std::string out; + out.reserve(1 + 8 * 8 + body.size() + 128); + out.push_back(static_cast(kEntryBlobVersion)); + put_field(out, std::to_string(status)); + put_field(out, etag); + put_field(out, last_modified); + put_field(out, content_type); + put_field(out, effective_url); + put_field(out, fetched_at_iso); + put_field(out, std::to_string(freshness_ttl)); + put_field(out, body); + return out; +} + +EntrySnapshot decode_entry(const std::string &blob) { EntrySnapshot snap; - cvc::state *e = entries.findDescendant(key); - if (!e) + if (blob.empty() || static_cast(blob[0]) != kEntryBlobVersion) return snap; - const boost::any d = e->data(); // by value — bind before casting a pointer into it - const std::string *body = boost::any_cast(&d); - if (!body) - return snap; // no committed body → miss + std::size_t pos = 1; + std::string status_ignored, + freshness_str; // status is stored for fidelity; the read path ignores it + if (!get_field(blob, pos, status_ignored) || !get_field(blob, pos, snap.etag) || + !get_field(blob, pos, snap.last_modified) || !get_field(blob, pos, snap.content_type) || + !get_field(blob, pos, snap.effective_url) || !get_field(blob, pos, snap.fetched_at) || + !get_field(blob, pos, freshness_str) || !get_field(blob, pos, snap.body)) + return snap; // malformed → miss + snap.freshness_ttl = freshness_str; snap.present = true; - snap.body = *body; - snap.etag = child_value(e, "etag"); - snap.last_modified = child_value(e, "last_modified"); - snap.content_type = child_value(e, "content_type"); - snap.effective_url = child_value(e, "effective_url"); - snap.fetched_at = child_value(e, "fetched_at"); - snap.freshness_ttl = child_value(e, "freshness_ttl"); return snap; } +// Read the entry at `key`. findDescendantShared PINS the node so a concurrent sweepExpired() can +// only unlink it (never free it under us); its single data() blob is then decoded — an atomic, +// parent-free read that needs no cache lock. Absent / non-blob / malformed all read as a miss. +EntrySnapshot read_entry(cvc::state &entries, const std::string &key) { + cvc::state::state_ptr e = entries.findDescendantShared(key); + if (!e) + return {}; + const boost::any d = e->data(); // by value — bind before casting a pointer into it + const std::string *blob = boost::any_cast(&d); + if (!blob) + return {}; + return decode_entry(*blob); +} + bool snapshot_fresh(const EntrySnapshot &s, const pt::ptime &now) { const long ttl = parse_long_or(s.freshness_ttl, 0); if (ttl <= 0 || s.fetched_at.empty()) @@ -301,30 +343,24 @@ bool snapshot_fresh(const EntrySnapshot &s, const pt::ptime &now) { return now < fetched + pt::seconds(ttl); } -// Write (create or replace) the full entry. Metadata children first, BODY LAST, then retention. -// `fetched_at` is the FRESHNESS anchor (receipt time minus the response Age). CALLER HOLDS -// cache_mutex, so a concurrent reader/sweeper cannot observe a torn entry or free the node -// mid-write. +// Write (create or replace) the entry as a SINGLE data() blob, then set retention. One atomic +// commit per entry node; single-flight already serializes writers of the same key, so no cache lock +// is needed. `fetched_at` is the FRESHNESS anchor (receipt time minus the response Age). void write_entry(cvc::state &entries, const std::string &key, long status, const std::string &etag, const std::string &last_modified, const std::string &content_type, const std::string &effective_url, const std::string &body, const pt::ptime &fetched_at, long freshness_ttl) { - cvc::state &e = entries(key); // create-or-get - e("status").value(std::to_string(status)); - e("etag").value(etag); - e("last_modified").value(last_modified); - e("content_type").value(content_type); - e("effective_url").value(effective_url); - e("fetched_at").value(pt::to_iso_string(fetched_at)); - e("freshness_ttl").value(std::to_string(freshness_ttl)); - e.data(boost::any(body)); // body LAST — the commit signal for a reader + cvc::state &e = entries(key); // create-or-get the entry node + e.data(boost::any(encode_entry(status, etag, last_modified, content_type, effective_url, + pt::to_iso_string(fetched_at), freshness_ttl, body))); const long retention = std::max(kRetentionSecs, freshness_ttl); e.expireAt(fetched_at + pt::seconds(retention)); } -// Force-expire + sweep the entry for `key` (cache invalidation). CALLER HOLDS cache_mutex. -void evict_locked(cvc::state &entries, const std::string &key, const pt::ptime &now) { - if (cvc::state *e = entries.findDescendant(key)) +// Force-expire + sweep the entry for `key` (cache invalidation). The pin keeps the node alive +// across expireAt; the sweep is safe against a concurrent pinned reader. +void evict(cvc::state &entries, const std::string &key, const pt::ptime &now) { + if (cvc::state::state_ptr e = entries.findDescendantShared(key)) e->expireAt(now - pt::seconds(1)); entries.sweepExpired(); } @@ -362,7 +398,7 @@ UriResult do_fetch_and_store(cvc::state &root, const std::string &key, const Uri req.headers.push_back("If-Modified-Since: " + stale.last_modified); } - const cvc::net::HttpResponse r = cvc::net::send(req); // NOT under cache_mutex + const cvc::net::HttpResponse r = cvc::net::send(req); // blocking send; holds no cache lock if (!r.ok) return {false, std::string(), std::string(), "ari: http fetch of '" + u.raw + "' failed: " + r.error}; @@ -373,8 +409,7 @@ UriResult do_fetch_and_store(cvc::state &root, const std::string &key, const Uri if (r.status == 304 && stale.present) { const CacheControl cc = parse_cache_control(r.headers); if (cc.no_store) { - std::lock_guard lk(cache_mutex()); - evict_locked(root(kCacheRoot), key, now); // must not retain; serve this body once + evict(root(kCacheRoot), key, now); // must not retain; serve this body once return {true, stale.body, stale.effective_url, std::string()}; } std::string etag = stale.etag, last_modified = stale.last_modified, @@ -386,11 +421,8 @@ UriResult do_fetch_and_store(cvc::state &root, const std::string &key, const Uri : cc.has_max_age ? std::max(0, cc.max_age) : parse_long_or(stale.freshness_ttl, kDefaultFreshnessSecs); const long age = parse_age(r.headers); - { - std::lock_guard lk(cache_mutex()); - write_entry(root(kCacheRoot), key, 200, etag, last_modified, content_type, - stale.effective_url, stale.body, now - pt::seconds(age), ttl); - } + write_entry(root(kCacheRoot), key, 200, etag, last_modified, content_type, stale.effective_url, + stale.body, now - pt::seconds(age), ttl); return {true, stale.body, stale.effective_url, std::string()}; } @@ -403,7 +435,6 @@ UriResult do_fetch_and_store(cvc::state &root, const std::string &key, const Uri find_header(r.headers, "Last-Modified", last_modified); find_header(r.headers, "Content-Type", content_type); const long age = parse_age(r.headers); - std::lock_guard lk(cache_mutex()); write_entry(root(kCacheRoot), key, 200, etag, last_modified, content_type, r.canonical_url, r.body, now - pt::seconds(age), derive_freshness_ttl(r.headers, cc)); } @@ -417,8 +448,8 @@ UriResult do_fetch_and_store(cvc::state &root, const std::string &key, const Uri // single-flight wrapper: the first caller for a key becomes the leader (does the fetch + store); // the rest wait and receive the leader's result. Coalesces only the network-bound paths (a fresh -// hit returns before this is called). cache_mutex is NEVER held here (the leader takes it inside -// do_fetch_and_store, after releasing flight_mutex). +// hit returns before this is called). This per-key serialization is the ONLY write coordination the +// cache needs — the entry write itself is a single atomic node op (see the CONCURRENCY note). UriResult fetch_and_cache(cvc::state &root, const std::string &key, const Uri &u, const EntrySnapshot &stale, const pt::ptime &now) { std::shared_ptr mine; @@ -474,10 +505,9 @@ UriResult cached_http_fetch(cvc::state &root, const Uri &u, const std::string & const pt::ptime now = pt::microsec_clock::universal_time(); EntrySnapshot snap; { - std::lock_guard lk(cache_mutex()); // sweep + snapshot atomically vs eviction cvc::state &entries = root(kCacheRoot); - entries.sweepExpired(); // access-time eviction (frees expired bodies) - snap = read_entry(entries, key); + entries.sweepExpired(); // access-time eviction (frees expired blobs) + snap = read_entry(entries, key); // read_entry PINS the node — safe against the sweep above } if (snap.present && snapshot_fresh(snap, now)) return {true, snap.body, snap.effective_url, std::string()}; // FRESH: no network @@ -489,10 +519,8 @@ UriResult cached_http_fetch(cvc::state &root, const Uri &u, const std::string & StoreResult invalidating_store(cvc::state &root, const Uri &u, const std::string &content, const std::string &base) { const StoreResult r = detail::http_store_uncached(u, content, base); - if (r.ok) { - std::lock_guard lk(cache_mutex()); - evict_locked(root(kCacheRoot), cache_key(u.raw), pt::microsec_clock::universal_time()); - } + if (r.ok) + evict(root(kCacheRoot), cache_key(u.raw), pt::microsec_clock::universal_time()); return r; } From ae8eaf50d9d98ac140764690870368be3cc19ef8 Mon Sep 17 00:00:00 2001 From: Joe Rivera Date: Tue, 29 Sep 2026 22:49:38 -0500 Subject: [PATCH 2/2] =?UTF-8?q?fix(state):=20pin=20the=20entry=20node=20in?= =?UTF-8?q?=20write=5Fentry=20(Phase=202=20review=20=E2=80=94=20write-vs-e?= =?UTF-8?q?vict=20UAF)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adversarial review of the cache_mutex removal found a HIGH UAF I introduced: with cache_mutex gone, write_entry still used entries(key) (operator(), a bare unpinned state&) and then mutated it via data()/expireAt, which lock the node and walk its _parent. A concurrent evict/sweepExpired (a same-key store invalidation — an independent URI handler no longer serialized against the fetch write) could free the node between entries(key) returning and the data()/expireAt calls → use-after-free (or a silently-lost write). read_entry and evict were already hardened with owning accessors; write_entry was the one mutation site left on a bare ref. Fix: write_entry now uses entries.sharedChild(key) (owning create-or-get), so the entry node is PINNED across data()/expireAt — a concurrent evict can only unlink it, never free it under the write. The entry's ancestors are the fixed, never-expiring sys.net.http_cache.entries path, so the entry pin alone suffices (no chain needed). Also refreshed the stale header comment to describe the single-blob storage. 12/12 HttpCache tests pass. --- src/cvc/ariadne/uri_http_cache.cpp | 25 ++++++++++++++++--------- 1 file changed, 16 insertions(+), 9 deletions(-) diff --git a/src/cvc/ariadne/uri_http_cache.cpp b/src/cvc/ariadne/uri_http_cache.cpp index 82b67703..8ca996bc 100644 --- a/src/cvc/ariadne/uri_http_cache.cpp +++ b/src/cvc/ariadne/uri_http_cache.cpp @@ -3,14 +3,17 @@ // It sits in front of the plain handler's transport + policy (cvc::net + uri_http_detail.h) and // adds a state-tree cache under `sys.net.http_cache.entries.`: // - key = cvc::sha256_hex(normalized URL) — a clean 64-hex path segment. -// - body = the entry node's data() channel (a raw std::string blob; state_exec reads it -// back as `bytes`, and it serves via state://…?data — the precursor bridge). -// - metadata = child value nodes (status/etag/last_modified/content_type/effective_url/ -// fetched_at/freshness_ttl) — replication-light strings. +// - entry = ONE opaque blob on the entry node's data(): a version byte + 8 +// length-prefixed +// fields (status/etag/last_modified/content_type/effective_url/fetched_at/ +// freshness_ttl/body). Packing everything into a single value makes a +// read/write one atomic node op (no torn multi-node entry) — see the +// CONCURRENCY note. // - freshness = SOFT: now < fetched_at + freshness_ttl → serve with NO network; else a // conditional GET revalidates (304 bump / 200 replace). // - retention = HARD: the node's expireAt (>= freshness) + sweepExpired() (swept on access) -// evicts and frees the body blob — one mechanism, the state tree's own expiry. +// evicts and frees the entry blob — one mechanism, the state tree's own +// expiry. // - single-flight = concurrent requests for one URL coalesce to a single transfer. // // CONCURRENCY: each entry is stored as ONE opaque blob on its node's data() (never as separate @@ -350,11 +353,15 @@ void write_entry(cvc::state &entries, const std::string &key, long status, const const std::string &last_modified, const std::string &content_type, const std::string &effective_url, const std::string &body, const pt::ptime &fetched_at, long freshness_ttl) { - cvc::state &e = entries(key); // create-or-get the entry node - e.data(boost::any(encode_entry(status, etag, last_modified, content_type, effective_url, - pt::to_iso_string(fetched_at), freshness_ttl, body))); + // sharedChild PINS the entry node (create-or-get), so a concurrent evict/sweepExpired (e.g. a + // store invalidation of the same key) can only unlink it, never free it under the data()/expireAt + // below — which lock the node and walk its _parent. The entry's ancestors are the fixed, + // never-expiring sys.net.http_cache.entries path, so pinning the entry alone is sufficient here. + cvc::state::state_ptr e = entries.sharedChild(key); + e->data(boost::any(encode_entry(status, etag, last_modified, content_type, effective_url, + pt::to_iso_string(fetched_at), freshness_ttl, body))); const long retention = std::max(kRetentionSecs, freshness_ttl); - e.expireAt(fetched_at + pt::seconds(retention)); + e->expireAt(fetched_at + pt::seconds(retention)); } // Force-expire + sweep the entry for `key` (cache invalidation). The pin keeps the node alive