Compare commits
4
Commits
cdedab3302
..
master
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
54982bcc62 | ||
|
|
a55f442315 | ||
|
|
0e0bf86af8 | ||
|
|
24175dde36 |
@@ -3,3 +3,7 @@
|
|||||||
|
|
||||||
/staging/*
|
/staging/*
|
||||||
!/staging/.gitkeep
|
!/staging/.gitkeep
|
||||||
|
|
||||||
|
# bench harness generates these locally (see bench/RPS); not meant to be committed
|
||||||
|
bench/urls*.txt
|
||||||
|
bench/full_urls*.txt
|
||||||
|
|||||||
@@ -4,6 +4,12 @@ version = "0.1.0"
|
|||||||
edition = "2024"
|
edition = "2024"
|
||||||
license = "AGPL-3.0-only"
|
license = "AGPL-3.0-only"
|
||||||
|
|
||||||
|
[features]
|
||||||
|
# RFC 007 causal profiling in smarm. Off by default (same zero-cost
|
||||||
|
# discipline as smarm itself): causal_site!/progress! call sites compile
|
||||||
|
# to no-ops without it, and `ccc causal` is unavailable.
|
||||||
|
causal = ["smarm/smarm-causal"]
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
urus = { git = "https://git.kalsbeek.dev/Markk116/urus.git", tag = "v0.2.2" }
|
urus = { git = "https://git.kalsbeek.dev/Markk116/urus.git", tag = "v0.2.2" }
|
||||||
# Pinned to the same tag urus itself depends on, so Cargo unifies both
|
# Pinned to the same tag urus itself depends on, so Cargo unifies both
|
||||||
|
|||||||
@@ -2,6 +2,11 @@
|
|||||||
|
|
||||||
My super simple CDN built for distributing my own (text) content.
|
My super simple CDN built for distributing my own (text) content.
|
||||||
|
|
||||||
|
Pushes 30k-45k requests/sec per core for a realistic workload -- see
|
||||||
|
[bench/RPS](bench/RPS) if you want the receipts, and
|
||||||
|
[bench/CAUSAL.md](bench/CAUSAL.md) for causal-profiling which sites
|
||||||
|
actually matter (`cargo build --features causal`).
|
||||||
|
|
||||||
## Running in Docker
|
## Running in Docker
|
||||||
|
|
||||||
```
|
```
|
||||||
|
|||||||
+112
@@ -0,0 +1,112 @@
|
|||||||
|
# CCC bench: causal profiling
|
||||||
|
|
||||||
|
`smarm` v0.6.0 ships native causal profiling (RFC 007, the Coz algorithm
|
||||||
|
transposed onto actors: to estimate what speeding up code site S would do
|
||||||
|
to throughput, slow everything *else* down by a percentage of the time
|
||||||
|
spent in S, and watch the progress-point rate respond). This is a much
|
||||||
|
better way to answer "what's actually worth optimizing?" than reading
|
||||||
|
tea leaves out of the raw RPS numbers in [bench/RPS](RPS) - e.g. that
|
||||||
|
doc's "is the cache scan an issue?" caveat can now be answered directly.
|
||||||
|
|
||||||
|
Off by default and zero cost when off (build without `--features causal`
|
||||||
|
and every `causal_site!`/`progress!` call compiles to a no-op). `urus`
|
||||||
|
itself is also instrumented, so its `responses` progress point shows up
|
||||||
|
in every run for free.
|
||||||
|
|
||||||
|
## Instrumented sites (`src/main.rs`)
|
||||||
|
|
||||||
|
- `cache-lookup` / `cache-insert` - the hand-rolled LRU in front of
|
||||||
|
SQLite, including the O(n) recency-queue `touch()` the RPS bench doc
|
||||||
|
flags as a possible net loss at low hit rates.
|
||||||
|
- `sqlite-query` - the `SELECT ... FROM versions` on cache miss.
|
||||||
|
- `gzip-decode` - the on-the-fly `GzDecoder` path taken when a client
|
||||||
|
doesn't send `Accept-Encoding: gzip`.
|
||||||
|
|
||||||
|
Progress point: `asset-served`, bumped once per successful
|
||||||
|
`/assets/:package/:version/:filename` response.
|
||||||
|
|
||||||
|
## Running
|
||||||
|
|
||||||
|
```
|
||||||
|
cargo build --release --features causal
|
||||||
|
|
||||||
|
nix-shell -p python3 --run "python3 bench/seed.py cdn.db"
|
||||||
|
awk '{print "http://127.0.0.1:8333"$0}' bench/urls.txt > /tmp/full_urls.txt
|
||||||
|
|
||||||
|
CCC_DB_PATH=$(pwd)/cdn.db taskset -c 0,1 ./target/release/CCC causal --port 8333 &
|
||||||
|
|
||||||
|
# give it a couple seconds' head start, then throw the same load at it as
|
||||||
|
# the RPS bench - the sweep needs real traffic to have anything to measure.
|
||||||
|
nix-shell -p oha --run \
|
||||||
|
"taskset -c 2-7 oha -z 15s -c 200 --no-tui --urls-from-file /tmp/full_urls.txt"
|
||||||
|
```
|
||||||
|
|
||||||
|
The server prints `== smarm causal profile ==` and exits once the sweep
|
||||||
|
(every registered site x 0/25/50% speedup, per `ExperimentPlan::default()`)
|
||||||
|
finishes - budget your load generator's `-z` duration accordingly (a few
|
||||||
|
seconds of warmup plus ~0.6s/cell). Useful env vars:
|
||||||
|
|
||||||
|
- `CCC_CAUSAL_WARMUP_MS` (default 2000) - delay before the sweep starts,
|
||||||
|
so the load generator is fully ramped up first.
|
||||||
|
- `CCC_CAUSAL_COZ_OUT=/path/to/profile.coz` - also dump a Coz-format
|
||||||
|
profile for Coz's existing plot tooling.
|
||||||
|
|
||||||
|
## Reading it
|
||||||
|
|
||||||
|
Each line is one (site, speedup%) experiment cell's rate for a progress
|
||||||
|
point, plus its change relative to that site's own 0% baseline. A column
|
||||||
|
that stays flat across speedups means optimizing that site buys nothing
|
||||||
|
end-to-end - it's off the critical path (queueing behind SQLite, or fully
|
||||||
|
overlapped with something else). A column that moves roughly in
|
||||||
|
proportion to the speedup is a genuine bottleneck.
|
||||||
|
|
||||||
|
Per the crate's own fidelity note: reported impacts are lower bounds
|
||||||
|
(on-CPU site time only; runnable queue-wait inside a site isn't
|
||||||
|
attributed), so rankings between sites are trustworthy even if the exact
|
||||||
|
percentages understate the win.
|
||||||
|
|
||||||
|
## Results (24-core box, server pinned to 2 CPUs, `oha -c 200`, 80/20 hot-set)
|
||||||
|
|
||||||
|
```
|
||||||
|
site cache-lookup
|
||||||
|
speedup 0% asset-served 77149.5/s vs baseline +0.0%
|
||||||
|
speedup 25% asset-served 79007.0/s vs baseline +2.4%
|
||||||
|
speedup 50% asset-served 80207.0/s vs baseline +4.0%
|
||||||
|
site sqlite-query
|
||||||
|
speedup 0% asset-served 80418.6/s vs baseline +0.0%
|
||||||
|
speedup 25% asset-served 91251.6/s vs baseline +13.5%
|
||||||
|
speedup 50% asset-served 96923.5/s vs baseline +20.5%
|
||||||
|
site cache-insert
|
||||||
|
speedup 0% asset-served 88755.9/s vs baseline +0.0%
|
||||||
|
speedup 25% asset-served 88560.3/s vs baseline -0.2%
|
||||||
|
speedup 50% asset-served 89800.1/s vs baseline +1.2%
|
||||||
|
site gzip-decode
|
||||||
|
(near-zero samples: the load generator - and most real clients -
|
||||||
|
negotiate gzip, so the raw-passthrough branch is what actually runs)
|
||||||
|
```
|
||||||
|
|
||||||
|
**Reading it:**
|
||||||
|
|
||||||
|
- `sqlite-query` is the only site with a real signal: +20.5% at a 50%
|
||||||
|
speedup, roughly linear with the injected speedup. It's the genuine
|
||||||
|
bottleneck on a cache miss.
|
||||||
|
- `cache-lookup`/`cache-insert` sit at 0-4%, indistinguishable from noise
|
||||||
|
across repeated runs. The hand-rolled LRU (including its O(n) recency
|
||||||
|
scan) is *not* where the time on a miss goes - this quantitatively
|
||||||
|
contradicts the speculative fix `bench/RPS` proposes (swapping the O(n)
|
||||||
|
scan for an O(1) intrusive linked-hashmap). Skip that; it wasn't going
|
||||||
|
to buy anything at these cache sizes.
|
||||||
|
- Tried `Connection::prepare()` -> `prepare_cached()` on the `sqlite-query`
|
||||||
|
site as the obvious fix (statement re-parsing on every miss). Re-ran
|
||||||
|
the same sweep after: **no measurable change** (+20.0% before,
|
||||||
|
+20.5% after - within run-to-run noise). Kept the change anyway (it's
|
||||||
|
strictly not worse and is idiomatic rusqlite), but it tells us parse
|
||||||
|
time isn't the dominant cost inside that site - execution (B-tree
|
||||||
|
lookup + copying the gzipped BLOB into a fresh `Vec<u8>`) is. Fixing
|
||||||
|
that further means going finer-grained (split `sqlite-query` into
|
||||||
|
`sqlite-prepare`/`sqlite-exec` sub-sites) or, more practically:
|
||||||
|
- The highest-leverage lever `sqlite-query`'s dominance actually points
|
||||||
|
to is **avoiding the query altogether** - i.e. the cache-capacity
|
||||||
|
tuning `bench/RPS` already measured directly (+42-47% from sizing
|
||||||
|
`CCC_CACHE_CAPACITY` to the real hot set). Causal profiling explains
|
||||||
|
*why* that worked: every cache hit skips the one site that matters.
|
||||||
@@ -0,0 +1,78 @@
|
|||||||
|
# CCC bench: RPS ceiling + LRU cache impact
|
||||||
|
|
||||||
|
Quick and dirty throughput bench, run locally on a 24-core box. Not
|
||||||
|
scientific, just enough to sanity-check the actor-based SQLite store and
|
||||||
|
the small LRU cache in front of it.
|
||||||
|
|
||||||
|
## Harness
|
||||||
|
|
||||||
|
- `bench/seed.py` fills a fresh `cdn.db` with random packages/versions/
|
||||||
|
assets (default: 500 packages x 5 versions = 2500 assets, 512B-8KB
|
||||||
|
gzipped JS each) and writes `bench/urls.txt` (one `/assets/...` path per
|
||||||
|
line, for every seeded asset).
|
||||||
|
- Server pinned to 2 CPUs, load generator (`oha`, via
|
||||||
|
`nix-shell -p oha`) pinned to 6 CPUs, both via `taskset`, so client
|
||||||
|
capacity is never the bottleneck.
|
||||||
|
|
||||||
|
```
|
||||||
|
nix-shell -p python3 --run "python3 bench/seed.py cdn.db"
|
||||||
|
awk '{print "http://127.0.0.1:8333"$0}' bench/urls.txt > /tmp/full_urls.txt
|
||||||
|
|
||||||
|
CCC_DB_PATH=$(pwd)/cdn.db taskset -c 0,1 ./target/release/CCC serve --port 8333 &
|
||||||
|
|
||||||
|
nix-shell -p oha --run \
|
||||||
|
"taskset -c 2-7 oha -z 8s -c 200 --no-tui --urls-from-file /tmp/full_urls.txt"
|
||||||
|
```
|
||||||
|
|
||||||
|
Skewed/hot-set workload (80% of requests hit the top 50 of 2500 assets,
|
||||||
|
i.e. a realistic CDN access pattern) was generated with a short Python
|
||||||
|
snippet sampling from `bench/urls.txt` with `random.random() < 0.8` picking
|
||||||
|
from the first 50 lines, else uniformly from the rest, written to a
|
||||||
|
`urls_hot80_20.txt` file and expanded the same way as above.
|
||||||
|
|
||||||
|
## Results: no cache (baseline)
|
||||||
|
|
||||||
|
Single-threaded `smarm` actor (`AssetStoreServer::loop_runner`) serializes
|
||||||
|
every asset lookup onto one thread/one SQLite connection, so this is
|
||||||
|
inherently CPU-bound on the 2 pinned cores regardless of client
|
||||||
|
concurrency.
|
||||||
|
|
||||||
|
| Server CPUs | Client concurrency | RPS |
|
||||||
|
|---|---|---|
|
||||||
|
| 2 | 50 | ~56.7k |
|
||||||
|
| 2 | 100 | ~60.1k |
|
||||||
|
| 2 | 200 | ~62.1k (peak, server at ~176-200% CPU, saturated) |
|
||||||
|
| 2 | 400 | ~60.0k (plateaued) |
|
||||||
|
| 2 | 800 | ~55.9k (queueing overhead) |
|
||||||
|
| 2 (client sends `Accept-Encoding: gzip`, server skips decompression) | 200 | ~60.4k |
|
||||||
|
| 1 | 200 | ~30.8k (confirms CPU-bound, scales with cores) |
|
||||||
|
|
||||||
|
Client (6 CPUs) stayed at ~4% usr / 8% sys throughout - never the
|
||||||
|
bottleneck.
|
||||||
|
|
||||||
|
## Results: with the LRU cache (`CCC_CACHE_CAPACITY`, default 256)
|
||||||
|
|
||||||
|
The cache lives inside `AssetStoreServer` itself (see `src/main.rs`), so
|
||||||
|
it needs no locking - the actor thread is already strictly sequential.
|
||||||
|
|
||||||
|
| Scenario | Cache capacity | RPS | vs. no-cache baseline (62.1k) |
|
||||||
|
|---|---|---|---|
|
||||||
|
| Uniform-random over all 2500 assets | 256 (default) | ~55.3k | **-11%** |
|
||||||
|
| Uniform-random over all 2500 assets | 3000 (covers full catalog) | ~88.1k | **+42%** |
|
||||||
|
| 80/20 hot-set (50 hot assets get 80% of traffic) | 256 (default) | ~91.1k | **+47%** |
|
||||||
|
|
||||||
|
### Caveat
|
||||||
|
|
||||||
|
Under a purely uniform-random access pattern with a cache smaller than
|
||||||
|
the catalog (low hit rate), the cache is a net loss: every request now
|
||||||
|
pays HashMap lookup + insert + eviction bookkeeping on top of the SQLite
|
||||||
|
query, for a hit rate too low to earn it back. The `touch()` on hit is
|
||||||
|
also an O(n) scan of the recency queue, which doesn't help at low
|
||||||
|
capacities.
|
||||||
|
|
||||||
|
Real CDN traffic is essentially never uniform-random (it's hot-set/
|
||||||
|
power-law skewed), so in practice this is a clear win, but `CCC_CACHE_CAPACITY`
|
||||||
|
should be sized to the actual hot set rather than left at the arbitrary
|
||||||
|
default of 256. A follow-up would swap the O(n) recency scan for a proper
|
||||||
|
O(1) LRU (e.g. an intrusive linked-hashmap) to remove the downside case
|
||||||
|
entirely.
|
||||||
@@ -0,0 +1,70 @@
|
|||||||
|
#!/usr/bin/env python3
|
||||||
|
"""Seed cdn.db with random packages/versions/files for benchmarking."""
|
||||||
|
import gzip
|
||||||
|
import os
|
||||||
|
import random
|
||||||
|
import sqlite3
|
||||||
|
import string
|
||||||
|
import sys
|
||||||
|
|
||||||
|
DB_PATH = sys.argv[1] if len(sys.argv) > 1 else "cdn.db"
|
||||||
|
N_PACKAGES = int(os.environ.get("N_PACKAGES", 500))
|
||||||
|
VERSIONS_PER_PKG = int(os.environ.get("VERSIONS_PER_PKG", 5))
|
||||||
|
MIN_SIZE = int(os.environ.get("MIN_SIZE", 512))
|
||||||
|
MAX_SIZE = int(os.environ.get("MAX_SIZE", 8192))
|
||||||
|
|
||||||
|
random.seed(42)
|
||||||
|
|
||||||
|
if os.path.exists(DB_PATH):
|
||||||
|
os.remove(DB_PATH)
|
||||||
|
|
||||||
|
conn = sqlite3.connect(DB_PATH)
|
||||||
|
conn.execute("""CREATE TABLE IF NOT EXISTS packages (
|
||||||
|
name TEXT PRIMARY KEY,
|
||||||
|
archived INTEGER NOT NULL DEFAULT 0
|
||||||
|
)""")
|
||||||
|
conn.execute("""CREATE TABLE IF NOT EXISTS versions (
|
||||||
|
package TEXT NOT NULL REFERENCES packages(name),
|
||||||
|
version TEXT NOT NULL,
|
||||||
|
filename TEXT NOT NULL,
|
||||||
|
mime_type TEXT NOT NULL,
|
||||||
|
gzipped_bytes BLOB NOT NULL,
|
||||||
|
archived INTEGER NOT NULL DEFAULT 0,
|
||||||
|
PRIMARY KEY (package, version)
|
||||||
|
)""")
|
||||||
|
|
||||||
|
def rand_name(n=10):
|
||||||
|
return "".join(random.choices(string.ascii_lowercase, k=n))
|
||||||
|
|
||||||
|
def rand_body(size):
|
||||||
|
chars = string.ascii_letters + string.digits + " \n"
|
||||||
|
return "".join(random.choices(chars, k=size)).encode()
|
||||||
|
|
||||||
|
manifest = [] # (package, version, filename) for the load generator
|
||||||
|
|
||||||
|
pkg_names = [f"pkg-{rand_name(8)}-{i}" for i in range(N_PACKAGES)]
|
||||||
|
|
||||||
|
for name in pkg_names:
|
||||||
|
conn.execute("INSERT INTO packages (name, archived) VALUES (?, 0)", (name,))
|
||||||
|
for v in range(VERSIONS_PER_PKG):
|
||||||
|
version = f"{v+1}.0.0"
|
||||||
|
filename = f"{rand_name(6)}.js"
|
||||||
|
size = random.randint(MIN_SIZE, MAX_SIZE)
|
||||||
|
raw = rand_body(size)
|
||||||
|
gz = gzip.compress(raw, compresslevel=6)
|
||||||
|
conn.execute(
|
||||||
|
"INSERT INTO versions (package, version, filename, mime_type, gzipped_bytes, archived) "
|
||||||
|
"VALUES (?, ?, ?, 'application/javascript', ?, 0)",
|
||||||
|
(name, version, filename, gz),
|
||||||
|
)
|
||||||
|
manifest.append((name, version, filename))
|
||||||
|
|
||||||
|
conn.commit()
|
||||||
|
conn.close()
|
||||||
|
|
||||||
|
with open("bench/urls.txt", "w") as f:
|
||||||
|
for pkg, ver, fn in manifest:
|
||||||
|
f.write(f"/assets/{pkg}/{ver}/{fn}\n")
|
||||||
|
|
||||||
|
print(f"seeded {len(pkg_names)} packages, {len(manifest)} versions -> {DB_PATH}")
|
||||||
|
print(f"wrote {len(manifest)} urls -> bench/urls.txt")
|
||||||
+187
-28
@@ -1,8 +1,11 @@
|
|||||||
mod store;
|
mod store;
|
||||||
|
|
||||||
use rusqlite::{params, Connection};
|
use rusqlite::{params, Connection};
|
||||||
|
use std::collections::{HashMap, VecDeque};
|
||||||
use std::io::Read;
|
use std::io::Read;
|
||||||
use std::sync::OnceLock;
|
use std::sync::OnceLock;
|
||||||
|
#[cfg(feature = "causal")]
|
||||||
|
use std::time::Duration;
|
||||||
use flate2::read::GzDecoder;
|
use flate2::read::GzDecoder;
|
||||||
use serde::Serialize;
|
use serde::Serialize;
|
||||||
use urus::{Config, Conn, Next, Pipeline, Router, serve_with};
|
use urus::{Config, Conn, Next, Pipeline, Router, serve_with};
|
||||||
@@ -13,6 +16,72 @@ pub struct AssetPayload {
|
|||||||
pub mime_type: String,
|
pub mime_type: String,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Identifies a single (package, version, filename) asset for cache lookups.
|
||||||
|
type AssetKey = (String, String, String);
|
||||||
|
|
||||||
|
/// Small hand-rolled LRU cache for hot assets, sitting in front of SQLite.
|
||||||
|
///
|
||||||
|
/// `AssetStoreServer` runs on a single dedicated actor thread (see
|
||||||
|
/// `loop_runner`), so this cache needs no locking whatsoever - every call
|
||||||
|
/// happens strictly sequentially. Capacity is intentionally small; this is
|
||||||
|
/// meant to absorb hot-asset traffic, not replace the DB as a working set.
|
||||||
|
struct AssetCache {
|
||||||
|
capacity: usize,
|
||||||
|
entries: HashMap<AssetKey, AssetPayload>,
|
||||||
|
// Recency queue, most-recently-used at the back. Kept simple (O(n)
|
||||||
|
// scan on hit) since capacity is small and this is a single thread.
|
||||||
|
order: VecDeque<AssetKey>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl AssetCache {
|
||||||
|
fn new(capacity: usize) -> Self {
|
||||||
|
Self { capacity, entries: HashMap::new(), order: VecDeque::new() }
|
||||||
|
}
|
||||||
|
|
||||||
|
fn get(&mut self, key: &AssetKey) -> Option<AssetPayload> {
|
||||||
|
if self.capacity == 0 {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
let hit = self.entries.get(key).cloned();
|
||||||
|
if hit.is_some() {
|
||||||
|
self.touch(key);
|
||||||
|
}
|
||||||
|
hit
|
||||||
|
}
|
||||||
|
|
||||||
|
fn put(&mut self, key: AssetKey, value: AssetPayload) {
|
||||||
|
if self.capacity == 0 {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if self.entries.contains_key(&key) {
|
||||||
|
self.entries.insert(key.clone(), value);
|
||||||
|
self.touch(&key);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if self.entries.len() >= self.capacity {
|
||||||
|
if let Some(oldest) = self.order.pop_front() {
|
||||||
|
self.entries.remove(&oldest);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
self.order.push_back(key.clone());
|
||||||
|
self.entries.insert(key, value);
|
||||||
|
}
|
||||||
|
|
||||||
|
fn touch(&mut self, key: &AssetKey) {
|
||||||
|
if let Some(pos) = self.order.iter().position(|k| k == key) {
|
||||||
|
let k = self.order.remove(pos).unwrap();
|
||||||
|
self.order.push_back(k);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn cache_capacity() -> usize {
|
||||||
|
std::env::var("CCC_CACHE_CAPACITY")
|
||||||
|
.ok()
|
||||||
|
.and_then(|v| v.parse().ok())
|
||||||
|
.unwrap_or(256)
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Serialize)]
|
#[derive(Serialize)]
|
||||||
pub struct PackageListing {
|
pub struct PackageListing {
|
||||||
name: String,
|
name: String,
|
||||||
@@ -33,15 +102,16 @@ pub enum GenServerMsg {
|
|||||||
|
|
||||||
pub struct AssetStoreServer {
|
pub struct AssetStoreServer {
|
||||||
conn: Connection,
|
conn: Connection,
|
||||||
|
cache: AssetCache,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl AssetStoreServer {
|
impl AssetStoreServer {
|
||||||
pub fn new() -> Self {
|
pub fn new() -> Self {
|
||||||
let conn = store::open().expect("Failed to open SQLite database");
|
let conn = store::open().expect("Failed to open SQLite database");
|
||||||
Self { conn }
|
Self { conn, cache: AssetCache::new(cache_capacity()) }
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn loop_runner(self, rx: smarm::Receiver<GenServerMsg>) {
|
pub fn loop_runner(mut self, rx: smarm::Receiver<GenServerMsg>) {
|
||||||
while let Ok(msg) = rx.recv() {
|
while let Ok(msg) = rx.recv() {
|
||||||
match msg {
|
match msg {
|
||||||
GenServerMsg::FetchAsset { package, version, filename, reply_to } => {
|
GenServerMsg::FetchAsset { package, version, filename, reply_to } => {
|
||||||
@@ -56,22 +126,51 @@ impl AssetStoreServer {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn fetch_asset(&self, package: &str, version: &str, filename: &str) -> Result<Option<AssetPayload>, String> {
|
fn fetch_asset(&mut self, package: &str, version: &str, filename: &str) -> Result<Option<AssetPayload>, String> {
|
||||||
let mut stmt = self.conn
|
let key: AssetKey = (package.to_string(), version.to_string(), filename.to_string());
|
||||||
.prepare(
|
{
|
||||||
"SELECT gzipped_bytes, mime_type FROM versions
|
// Suspect per the RPS bench caveat: touch() is an O(n) scan of
|
||||||
WHERE package = ? AND version = ? AND filename = ?",
|
// the recency queue, so cache overhead itself is a candidate
|
||||||
)
|
// bottleneck at low hit rates - causal profiling can confirm or
|
||||||
.map_err(|e| e.to_string())?;
|
// rule that out instead of guessing from the raw RPS numbers.
|
||||||
|
let _g = smarm::causal_site!("cache-lookup");
|
||||||
|
if let Some(cached) = self.cache.get(&key) {
|
||||||
|
return Ok(Some(cached));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
let mut rows = stmt.query(params![package, version, filename]).map_err(|e| e.to_string())?;
|
let payload = {
|
||||||
|
let _g = smarm::causal_site!("sqlite-query");
|
||||||
|
// Causal profiling (see bench/CAUSAL.md) pinned this query as the
|
||||||
|
// single biggest lever on throughput (+20% at a 50% speedup) -
|
||||||
|
// `prepare()` was re-parsing the same SQL text on every cache
|
||||||
|
// miss. `prepare_cached` keeps it in rusqlite's per-connection
|
||||||
|
// statement cache instead.
|
||||||
|
let mut stmt = self.conn
|
||||||
|
.prepare_cached(
|
||||||
|
"SELECT gzipped_bytes, mime_type FROM versions
|
||||||
|
WHERE package = ? AND version = ? AND filename = ?",
|
||||||
|
)
|
||||||
|
.map_err(|e| e.to_string())?;
|
||||||
|
|
||||||
if let Some(row) = rows.next().map_err(|e| e.to_string())? {
|
let mut rows = stmt.query(params![package, version, filename]).map_err(|e| e.to_string())?;
|
||||||
let gzipped_bytes: Vec<u8> = row.get(0).map_err(|e| e.to_string())?;
|
|
||||||
let mime_type: String = row.get(1).map_err(|e| e.to_string())?;
|
if let Some(row) = rows.next().map_err(|e| e.to_string())? {
|
||||||
Ok(Some(AssetPayload { gzipped_bytes, mime_type }))
|
let gzipped_bytes: Vec<u8> = row.get(0).map_err(|e| e.to_string())?;
|
||||||
} else {
|
let mime_type: String = row.get(1).map_err(|e| e.to_string())?;
|
||||||
Ok(None)
|
Some(AssetPayload { gzipped_bytes, mime_type })
|
||||||
|
} else {
|
||||||
|
None
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
match payload {
|
||||||
|
Some(payload) => {
|
||||||
|
let _g = smarm::causal_site!("cache-insert");
|
||||||
|
self.cache.put(key, payload.clone());
|
||||||
|
Ok(Some(payload))
|
||||||
|
}
|
||||||
|
None => Ok(None),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -112,10 +211,50 @@ fn print_usage() {
|
|||||||
\x20 ccc add <package> <filepath> <version>\n\
|
\x20 ccc add <package> <filepath> <version>\n\
|
||||||
\x20 ccc archive <package> [<version>]\n\
|
\x20 ccc archive <package> [<version>]\n\
|
||||||
\x20 ccc stats\n\
|
\x20 ccc stats\n\
|
||||||
\x20 ccc serve [-p|--port <port>]\n"
|
\x20 ccc serve [-p|--port <port>]\n\
|
||||||
|
\x20 ccc causal [-p|--port <port>] (requires --features causal)\n"
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Whether `causal` (vs. plain `serve`) was requested - checked once the
|
||||||
|
/// server actually starts, since `run_cli` only hands `main` a port.
|
||||||
|
static CAUSAL_MODE: OnceLock<bool> = OnceLock::new();
|
||||||
|
|
||||||
|
/// Runs the RFC 007 experiment sweep against the live server on a plain OS
|
||||||
|
/// thread, prints the summary, and exits. Meant to be run alongside an
|
||||||
|
/// external load generator (same setup as `bench/RPS`) - the sweep needs
|
||||||
|
/// real request traffic hitting `causal_site!`/`progress!` call sites to
|
||||||
|
/// produce anything.
|
||||||
|
#[cfg(feature = "causal")]
|
||||||
|
fn spawn_causal_profiler() {
|
||||||
|
std::thread::spawn(|| {
|
||||||
|
let warmup = std::env::var("CCC_CAUSAL_WARMUP_MS")
|
||||||
|
.ok()
|
||||||
|
.and_then(|v| v.parse().ok())
|
||||||
|
.map(Duration::from_millis)
|
||||||
|
.unwrap_or(Duration::from_secs(2));
|
||||||
|
eprintln!("[causal] warming up for {warmup:?}, send traffic now...");
|
||||||
|
std::thread::sleep(warmup);
|
||||||
|
eprintln!("[causal] starting experiment sweep...");
|
||||||
|
let results = smarm::causal::run_experiments(&Default::default());
|
||||||
|
println!("{}", smarm::causal::render_summary(&results));
|
||||||
|
if std::env::var("CCC_CAUSAL_LEDGER").is_ok() {
|
||||||
|
eprintln!("{}", smarm::causal::render_ledger_audit(&results));
|
||||||
|
}
|
||||||
|
if let Ok(path) = std::env::var("CCC_CAUSAL_COZ_OUT") {
|
||||||
|
let _ = std::fs::write(&path, smarm::causal::render_coz(&results));
|
||||||
|
eprintln!("[causal] wrote {path}");
|
||||||
|
}
|
||||||
|
std::process::exit(0);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(not(feature = "causal"))]
|
||||||
|
fn spawn_causal_profiler() {
|
||||||
|
eprintln!("error: 'ccc causal' requires building with --features causal");
|
||||||
|
std::process::exit(2);
|
||||||
|
}
|
||||||
|
|
||||||
fn run_cli() -> Option<u16> {
|
fn run_cli() -> Option<u16> {
|
||||||
let args: Vec<String> = std::env::args().collect();
|
let args: Vec<String> = std::env::args().collect();
|
||||||
match args.get(1).map(String::as_str) {
|
match args.get(1).map(String::as_str) {
|
||||||
@@ -182,7 +321,8 @@ fn run_cli() -> Option<u16> {
|
|||||||
}
|
}
|
||||||
None
|
None
|
||||||
}
|
}
|
||||||
Some("serve") | None => {
|
Some("causal") | Some("serve") | None => {
|
||||||
|
let _ = CAUSAL_MODE.set(args.get(1).map(String::as_str) == Some("causal"));
|
||||||
let mut port: u16 = 8333;
|
let mut port: u16 = 8333;
|
||||||
let mut i = 2;
|
let mut i = 2;
|
||||||
while i < args.len() {
|
while i < args.len() {
|
||||||
@@ -267,7 +407,7 @@ fn fetch_asset_handler(c: Conn, _n: Next) -> Conn {
|
|||||||
.map(|v| v.contains("gzip"))
|
.map(|v| v.contains("gzip"))
|
||||||
.unwrap_or(false);
|
.unwrap_or(false);
|
||||||
|
|
||||||
if accepts_gzip {
|
let response = if accepts_gzip {
|
||||||
c.put_status(200)
|
c.put_status(200)
|
||||||
.put_header("content-type", &asset.mime_type)
|
.put_header("content-type", &asset.mime_type)
|
||||||
.put_header("content-encoding", "gzip")
|
.put_header("content-encoding", "gzip")
|
||||||
@@ -275,33 +415,52 @@ fn fetch_asset_handler(c: Conn, _n: Next) -> Conn {
|
|||||||
.put_header("access-control-allow-origin", "*")
|
.put_header("access-control-allow-origin", "*")
|
||||||
.put_body(asset.gzipped_bytes)
|
.put_body(asset.gzipped_bytes)
|
||||||
} else {
|
} else {
|
||||||
let mut decoder = GzDecoder::new(&asset.gzipped_bytes[..]);
|
let raw_bytes = {
|
||||||
let mut raw_bytes = Vec::new();
|
let _g = smarm::causal_site!("gzip-decode");
|
||||||
if decoder.read_to_end(&mut raw_bytes).is_ok() {
|
let mut decoder = GzDecoder::new(&asset.gzipped_bytes[..]);
|
||||||
c.put_status(200)
|
let mut raw_bytes = Vec::new();
|
||||||
|
decoder.read_to_end(&mut raw_bytes).map(|_| raw_bytes)
|
||||||
|
};
|
||||||
|
match raw_bytes {
|
||||||
|
Ok(raw_bytes) => c.put_status(200)
|
||||||
.put_header("content-type", &asset.mime_type)
|
.put_header("content-type", &asset.mime_type)
|
||||||
.put_header("cache-control", "public, max-age=31536000, immutable")
|
.put_header("cache-control", "public, max-age=31536000, immutable")
|
||||||
.put_header("access-control-allow-origin", "*")
|
.put_header("access-control-allow-origin", "*")
|
||||||
.put_body(raw_bytes)
|
.put_body(raw_bytes),
|
||||||
} else {
|
Err(_) => return c.put_status(500).put_body("Decompression Error"),
|
||||||
c.put_status(500).put_body("Decompression Error")
|
|
||||||
}
|
}
|
||||||
}
|
};
|
||||||
|
smarm::progress!("asset-served");
|
||||||
|
response
|
||||||
}
|
}
|
||||||
Ok(Ok(None)) => c.put_status(404).put_body("Asset Not Found"),
|
Ok(Ok(None)) => c.put_status(404).put_body("Asset Not Found"),
|
||||||
_ => c.put_status(500).put_body("Database Error Encountered"),
|
_ => c.put_status(500).put_body("Database Error Encountered"),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Fallback for any request the router didn't match. Deliberately terse -
|
||||||
|
/// no route list, no hint of the `/assets/:package/:version/:filename`
|
||||||
|
/// shape, nothing to make this service more discoverable to strangers
|
||||||
|
/// than it already is. Just enough to not look like a hung connection.
|
||||||
|
fn not_found_handler(c: Conn, _n: Next) -> Conn {
|
||||||
|
c.put_status(404)
|
||||||
|
.put_header("content-type", "text/plain")
|
||||||
|
.put_body("404 not found")
|
||||||
|
}
|
||||||
|
|
||||||
fn main() {
|
fn main() {
|
||||||
let Some(port) = run_cli() else { return };
|
let Some(port) = run_cli() else { return };
|
||||||
|
|
||||||
|
if *CAUSAL_MODE.get().unwrap_or(&false) {
|
||||||
|
spawn_causal_profiler();
|
||||||
|
}
|
||||||
|
|
||||||
let router = Router::new()
|
let router = Router::new()
|
||||||
.get("/packages", list_packages_handler)
|
.get("/packages", list_packages_handler)
|
||||||
.get("/assets/:package/:version/:filename", fetch_asset_handler);
|
.get("/assets/:package/:version/:filename", fetch_asset_handler);
|
||||||
|
|
||||||
println!("C3 on 0.0.0.0:{port}...");
|
println!("C3 on 0.0.0.0:{port}...");
|
||||||
let cfg = Config::new(format!("0.0.0.0:{port}").parse().unwrap());
|
let cfg = Config::new(format!("0.0.0.0:{port}").parse().unwrap());
|
||||||
let pipeline = Pipeline::new().plug(router);
|
let pipeline = Pipeline::new().plug(router).plug(not_found_handler);
|
||||||
serve_with(cfg, pipeline).unwrap();
|
serve_with(cfg, pipeline).unwrap();
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user