Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
193 changes: 114 additions & 79 deletions src/cvc/ariadne/uri_http_cache.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -3,25 +3,28 @@
// 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>`:
// - 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: 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
Expand All @@ -35,6 +38,7 @@
#include <cctype>
#include <condition_variable>
#include <cstddef>
#include <cstdint>
#include <cvc/ariadne/uri.h>
#include <cvc/ariadne/uri_http.h>
#include <cvc/ariadne/uri_http_cache.h>
Expand Down Expand Up @@ -62,12 +66,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-
Expand Down Expand Up @@ -244,10 +242,9 @@ long derive_freshness_ttl(const std::vector<std::string> &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;
Expand All @@ -259,35 +256,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<char>((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<std::uint64_t>(static_cast<unsigned char>(in[pos + i])) << (8 * i);
pos += 8;
if (n > in.size() - pos)
return false;
out.assign(in, pos, static_cast<std::size_t>(n));
pos += static_cast<std::size_t>(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<char>(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<unsigned char>(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<std::string>(&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<std::string>(&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())
Expand All @@ -301,30 +346,28 @@ 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
// 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). 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();
}
Expand Down Expand Up @@ -362,7 +405,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};
Expand All @@ -373,8 +416,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<std::mutex> 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,
Expand All @@ -386,11 +428,8 @@ UriResult do_fetch_and_store(cvc::state &root, const std::string &key, const Uri
: cc.has_max_age ? std::max<long>(0, cc.max_age)
: parse_long_or(stale.freshness_ttl, kDefaultFreshnessSecs);
const long age = parse_age(r.headers);
{
std::lock_guard<std::mutex> 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()};
}

Expand All @@ -403,7 +442,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<std::mutex> 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));
}
Expand All @@ -417,8 +455,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<Inflight> mine;
Expand Down Expand Up @@ -474,10 +512,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<std::mutex> 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
Expand All @@ -489,10 +526,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<std::mutex> 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;
}

Expand Down
Loading