Compare commits

..
4 Commits
Author SHA1 Message Date
Markk116 a55f442315 Add causal profiling (RFC 007) behind --features causal
Instrument the four suspect regions from bench/RPS: cache-lookup,
cache-insert, sqlite-query, and gzip-decode, plus an asset-served
progress point. New `ccc causal` subcommand runs the sweep against
live traffic and prints a summary (optionally a .coz file and a
ledger audit).

Ran it under the same 80/20 hot-set workload as bench/RPS - see
bench/CAUSAL.md for the methodology and results. Findings:

- sqlite-query is the real bottleneck on a cache miss (+20.5% at a
  50% speedup, roughly linear).
- cache-lookup/cache-insert are noise-level (0-4%) - the O(n)
  recency-scan touch() bench/RPS flagged as a possible follow-up is
  not actually costing anything, so that's off the table.
- Switched prepare() -> prepare_cached() on the query as the obvious
  fix; re-measured and it made no real difference (+20.0% -> +20.5%,
  within noise). Kept it anyway (strictly not worse), but it shows
  execution cost (B-tree lookup + BLOB copy) dominates over parse
  cost in that site.
- gzip-decode barely gets exercised since real clients (and oha)
  negotiate gzip - not worth optimizing further.
- Net conclusion: the existing CCC_CACHE_CAPACITY tuning from
  bench/RPS (+42-47% RPS) is the correct lever, and causal profiling
  explains why - every cache hit skips the one site that matters.

Zero cost when the feature is off: causal_site!/progress! compile to
no-ops without smarm-causal.
2026-08-08 23:35:29 +02:00
Markk116 0e0bf86af8 README: add headline RPS claim, point to bench/RPS for methodology
30k-45k req/s/core is the measured range in bench/RPS (cold-SQLite
floor to LRU-cached hot-set ceiling). No inline caveats here on
purpose -- anyone who wants the fine print goes and reads the harness.
2026-08-08 23:01:46 +02:00
Markk116 24175dde36 Add small in-actor LRU cache for hot assets
AssetStoreServer runs on a single dedicated smarm actor thread, so a
plain HashMap+VecDeque LRU in front of the SQLite lookup needs no
locking. Capacity configurable via CCC_CACHE_CAPACITY (default 256).

Bench (bench/RPS): +42-47% RPS under a cache-sized/hot-set workload,
but a small net loss under adversarial uniform-random access with a
cache smaller than the catalog. Real traffic is hot-set skewed, so
net win in practice; capacity should be tuned to the expected hot set.
2026-08-08 22:59:57 +02:00
Markk116 cdedab3302 Add test suite, stats command, and fix runtime stack-overflow crash
- Bump urus to v0.2.2 (v0.2.1 -> v0.2.2) and smarm to v0.6.0, pulling in
  a urus fix (pushed alongside this commit) that spawns per-connection
  actors with a 256 KiB stack via smarm's RFC 019 SpawnOpts instead of
  the runtime's bare 64 KiB default. Without it, any request that hit
  fetch_asset_handler's in-handler gzip decompression (i.e. any client
  not sending Accept-Encoding: gzip) blew the actor's guard page and the
  connection died with no response - reproduced 5/5 runs before the fix,
  0/5 after.

- Add `ccc stats`: active/archived package and version counts, total
  gzipped bytes stored, and the resolved db path. Useful for a quick
  sanity check before/after a deploy.

- Add a real test suite, which is what caught the crash above:
    - src/store.rs unit tests: schema init/idempotency, create/add/
      archive success and error paths, gzip round-trip, upsert
      semantics, stats aggregation.
    - tests/cli.rs: black-box tests against the compiled `ccc` binary
      covering usage/exit codes and the full create/add/archive/stats
      lifecycle.
    - tests/server.rs: boots `ccc serve` as a real subprocess and drives
      it over raw TCP (no HTTP client dependency) - boot-without-
      crashing, asset serving both compressed and decompressed, 404s,
      package listing incl. archived-package hiding, and repeated
      sequential requests against one long-lived process.

27/27 tests passing (12 unit + 9 CLI + 6 server).
2026-08-08 22:42:37 +02:00
11 changed files with 1200 additions and 33 deletions
+4
View File
@@ -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
Generated
+5 -4
View File
@@ -378,9 +378,10 @@ checksum = "8ed6a63f02c8539c91a8685a86f4099661ba3da017932f6ebbea6de3f0fa7c90"
[[package]] [[package]]
name = "smarm" name = "smarm"
version = "0.5.0" version = "0.6.0"
source = "git+https://git.kalsbeek.dev/Markk116/smarm.git?tag=v0.5.0#a03a7ca01ef1ecaeb8aabc00bebe7ddef4f3879a" source = "git+https://git.kalsbeek.dev/Markk116/smarm.git?tag=v0.6.0#301e3463e3abb11ee25b94ee2ec31c497fb4e2b4"
dependencies = [ dependencies = [
"cc",
"libc", "libc",
"loom", "loom",
] ]
@@ -473,8 +474,8 @@ checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75"
[[package]] [[package]]
name = "urus" name = "urus"
version = "0.2.0" version = "0.2.2"
source = "git+https://git.kalsbeek.dev/Markk116/urus.git#b77448191ec33de1662305f3035473d64f8cd34f" source = "git+https://git.kalsbeek.dev/Markk116/urus.git?tag=v0.2.2#535f7bcc687bcdf4ef5cf366177ac166aa8466a1"
dependencies = [ dependencies = [
"httparse", "httparse",
"libc", "libc",
+8 -2
View File
@@ -4,11 +4,17 @@ 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" } 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
# into a single source id/copy of smarm (avoids duplicate #[global_allocator]). # into a single source id/copy of smarm (avoids duplicate #[global_allocator]).
smarm = { git = "https://git.kalsbeek.dev/Markk116/smarm.git", tag = "v0.5.0" } smarm = { git = "https://git.kalsbeek.dev/Markk116/smarm.git", tag = "v0.6.0" }
rusqlite = { version = "0.31", features = ["bundled"] } rusqlite = { version = "0.31", features = ["bundled"] }
flate2 = "1.0" flate2 = "1.0"
serde = { version = "1", features = ["derive"] } serde = { version = "1", features = ["derive"] }
+5
View File
@@ -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
View File
@@ -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.
+78
View File
@@ -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.
+70
View File
@@ -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")
+192 -27
View File
@@ -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),
} }
} }
@@ -111,10 +210,51 @@ fn print_usage() {
\x20 ccc create <package>\n\ \x20 ccc create <package>\n\
\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 serve [-p|--port <port>]\n" \x20 ccc stats\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) {
@@ -148,6 +288,21 @@ fn run_cli() -> Option<u16> {
} }
None None
} }
Some("stats") => {
match store::stats() {
Ok(s) => {
println!("packages: {} active, {} archived", s.packages_active, s.packages_archived);
println!("versions: {} active, {} archived", s.versions_active, s.versions_archived);
println!("stored: {} bytes (gzipped)", s.total_gzipped_bytes);
println!("db path: {}", store::db_path().display());
}
Err(e) => {
eprintln!("error: {e}");
std::process::exit(1);
}
}
None
}
Some("archive") => { Some("archive") => {
let Some(package) = args.get(2) else { let Some(package) = args.get(2) else {
print_usage(); print_usage();
@@ -166,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() {
@@ -251,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")
@@ -259,18 +415,23 @@ 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"),
@@ -280,6 +441,10 @@ fn fetch_asset_handler(c: Conn, _n: Next) -> Conn {
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);
+304
View File
@@ -127,6 +127,47 @@ pub fn archive(package: &str, version: Option<&str>) -> io::Result<()> {
Ok(()) Ok(())
} }
#[derive(Debug, PartialEq, Eq)]
pub struct Stats {
pub packages_active: i64,
pub packages_archived: i64,
pub versions_active: i64,
pub versions_archived: i64,
pub total_gzipped_bytes: i64,
}
pub fn stats() -> io::Result<Stats> {
let conn = open().map_err(io_err)?;
let packages_active: i64 = conn
.query_row("SELECT COUNT(*) FROM packages WHERE archived = 0", [], |r| r.get(0))
.map_err(io_err)?;
let packages_archived: i64 = conn
.query_row("SELECT COUNT(*) FROM packages WHERE archived = 1", [], |r| r.get(0))
.map_err(io_err)?;
let versions_active: i64 = conn
.query_row("SELECT COUNT(*) FROM versions WHERE archived = 0", [], |r| r.get(0))
.map_err(io_err)?;
let versions_archived: i64 = conn
.query_row("SELECT COUNT(*) FROM versions WHERE archived = 1", [], |r| r.get(0))
.map_err(io_err)?;
let total_gzipped_bytes: i64 = conn
.query_row(
"SELECT COALESCE(SUM(LENGTH(gzipped_bytes)), 0) FROM versions",
[],
|r| r.get(0),
)
.map_err(io_err)?;
Ok(Stats {
packages_active,
packages_archived,
versions_active,
versions_archived,
total_gzipped_bytes,
})
}
fn guess_mime(path: &std::path::Path) -> &'static str { fn guess_mime(path: &std::path::Path) -> &'static str {
match path.extension().and_then(|e| e.to_str()).unwrap_or("") { match path.extension().and_then(|e| e.to_str()).unwrap_or("") {
"js" | "mjs" => "application/javascript", "js" | "mjs" => "application/javascript",
@@ -144,3 +185,266 @@ fn guess_mime(path: &std::path::Path) -> &'static str {
_ => "application/octet-stream", _ => "application/octet-stream",
} }
} }
#[cfg(test)]
mod tests {
use super::*;
use std::io::Write;
use std::sync::Mutex;
// store::open()/db_path() read the process-wide CCC_DB_PATH env var, so
// tests that touch it must not run concurrently on separate threads.
static ENV_LOCK: Mutex<()> = Mutex::new(());
struct TempDb {
path: PathBuf,
_guard: std::sync::MutexGuard<'static, ()>,
}
impl TempDb {
fn new(tag: &str) -> Self {
let guard = ENV_LOCK.lock().unwrap_or_else(|e| e.into_inner());
let path = std::env::temp_dir().join(format!(
"ccc-test-{tag}-{}-{:?}.db",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let _ = std::fs::remove_file(&path);
unsafe { std::env::set_var("CCC_DB_PATH", &path) };
Self { path, _guard: guard }
}
}
impl Drop for TempDb {
fn drop(&mut self) {
unsafe { std::env::remove_var("CCC_DB_PATH") };
let _ = std::fs::remove_file(&self.path);
}
}
fn write_temp_file(name: &str, contents: &[u8]) -> PathBuf {
// Keep `name` (with its extension) as the trailing path component so
// add_version()'s guess_mime()/filename logic sees the real extension.
let dir = std::env::temp_dir().join(format!(
"ccc-src-{}-{:?}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join(name);
let mut f = std::fs::File::create(&path).unwrap();
f.write_all(contents).unwrap();
path
}
#[test]
fn db_path_defaults_to_cdn_db_when_unset() {
let _guard = ENV_LOCK.lock().unwrap_or_else(|e| e.into_inner());
unsafe { std::env::remove_var("CCC_DB_PATH") };
assert_eq!(db_path(), PathBuf::from("cdn.db"));
}
#[test]
fn open_creates_schema_and_is_idempotent() {
let db = TempDb::new("schema");
let conn = open().expect("first open should create schema");
drop(conn);
// Re-opening an existing DB must not fail on `CREATE TABLE IF NOT EXISTS`.
let conn2 = open().expect("second open should reuse schema");
let count: i64 = conn2
.query_row("SELECT COUNT(*) FROM packages", [], |r| r.get(0))
.unwrap();
assert_eq!(count, 0);
let _ = db;
}
#[test]
fn create_package_then_duplicate_errors() {
let _db = TempDb::new("create-dup");
create_package("demo").expect("create should succeed");
let err = create_package("demo").expect_err("duplicate create should fail");
assert_eq!(err.kind(), io::ErrorKind::AlreadyExists);
}
#[test]
fn add_version_without_package_errors() {
let _db = TempDb::new("add-missing-pkg");
let src = write_temp_file("nope.txt", b"hello");
let err = add_version("ghost", src.to_str().unwrap(), "1.0.0")
.expect_err("adding to a nonexistent package should fail");
assert_eq!(err.kind(), io::ErrorKind::NotFound);
let _ = std::fs::remove_file(&src);
}
#[test]
fn add_version_round_trips_gzip_and_mime() {
let _db = TempDb::new("add-roundtrip");
create_package("demo").unwrap();
let src = write_temp_file("app.js", b"console.log('hi');");
add_version("demo", src.to_str().unwrap(), "1.0.0").expect("add should succeed");
let conn = open().unwrap();
let (filename, mime, blob): (String, String, Vec<u8>) = conn
.query_row(
"SELECT filename, mime_type, gzipped_bytes FROM versions WHERE package = ? AND version = ?",
params!["demo", "1.0.0"],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.unwrap();
assert_eq!(mime, "application/javascript");
assert!(filename.ends_with("app.js") || filename.contains("app.js"));
use std::io::Read;
let mut decoder = flate2::read::GzDecoder::new(&blob[..]);
let mut decompressed = Vec::new();
decoder.read_to_end(&mut decompressed).unwrap();
assert_eq!(decompressed, b"console.log('hi');");
let _ = std::fs::remove_file(&src);
}
#[test]
fn add_version_upsert_replaces_existing_version() {
let _db = TempDb::new("add-upsert");
create_package("demo").unwrap();
let src1 = write_temp_file("a.txt", b"first");
let src2 = write_temp_file("b.css", b"body{color:red}");
add_version("demo", src1.to_str().unwrap(), "1.0.0").unwrap();
add_version("demo", src2.to_str().unwrap(), "1.0.0").unwrap();
let conn = open().unwrap();
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM versions WHERE package = ? AND version = ?",
params!["demo", "1.0.0"],
|r| r.get(0),
)
.unwrap();
assert_eq!(count, 1, "upsert should replace, not duplicate, the row");
let mime: String = conn
.query_row(
"SELECT mime_type FROM versions WHERE package = ? AND version = ?",
params!["demo", "1.0.0"],
|r| r.get(0),
)
.unwrap();
assert_eq!(mime, "text/css");
let _ = std::fs::remove_file(&src1);
let _ = std::fs::remove_file(&src2);
}
#[test]
fn stats_reports_active_archived_counts_and_bytes() {
let _db = TempDb::new("stats");
create_package("demo").unwrap();
create_package("other").unwrap();
let src1 = write_temp_file("a.txt", b"hello");
let src2 = write_temp_file("b.txt", b"world!!");
add_version("demo", src1.to_str().unwrap(), "1.0.0").unwrap();
add_version("demo", src2.to_str().unwrap(), "2.0.0").unwrap();
archive("demo", Some("1.0.0")).unwrap();
archive("other", None).unwrap();
let s = stats().expect("stats should succeed");
assert_eq!(s.packages_active, 1);
assert_eq!(s.packages_archived, 1);
assert_eq!(s.versions_active, 1);
assert_eq!(s.versions_archived, 1);
assert!(s.total_gzipped_bytes > 0, "gzipped bytes should be non-zero");
let _ = std::fs::remove_file(&src1);
let _ = std::fs::remove_file(&src2);
}
#[test]
fn stats_on_empty_db_is_all_zero() {
let _db = TempDb::new("stats-empty");
let s = stats().expect("stats should succeed on empty db");
assert_eq!(
s,
Stats {
packages_active: 0,
packages_archived: 0,
versions_active: 0,
versions_archived: 0,
total_gzipped_bytes: 0,
}
);
}
#[test]
fn archive_unknown_package_errors() {
let _db = TempDb::new("archive-missing");
let err = archive("ghost", None).expect_err("archiving unknown package should fail");
assert_eq!(err.kind(), io::ErrorKind::NotFound);
}
#[test]
fn archive_unknown_version_errors() {
let _db = TempDb::new("archive-missing-version");
create_package("demo").unwrap();
let err = archive("demo", Some("9.9.9"))
.expect_err("archiving unknown version should fail");
assert_eq!(err.kind(), io::ErrorKind::NotFound);
}
#[test]
fn archive_whole_package_hides_it_but_keeps_versions_row() {
let _db = TempDb::new("archive-package");
create_package("demo").unwrap();
let src = write_temp_file("a.txt", b"x");
add_version("demo", src.to_str().unwrap(), "1.0.0").unwrap();
archive("demo", None).expect("archive whole package should succeed");
let conn = open().unwrap();
let archived: i64 = conn
.query_row("SELECT archived FROM packages WHERE name = ?", params!["demo"], |r| r.get(0))
.unwrap();
assert_eq!(archived, 1);
let _ = std::fs::remove_file(&src);
}
#[test]
fn archive_single_version_only_affects_that_version() {
let _db = TempDb::new("archive-version");
create_package("demo").unwrap();
let src1 = write_temp_file("a.txt", b"x");
let src2 = write_temp_file("c.txt", b"y");
add_version("demo", src1.to_str().unwrap(), "1.0.0").unwrap();
add_version("demo", src2.to_str().unwrap(), "2.0.0").unwrap();
archive("demo", Some("1.0.0")).unwrap();
let conn = open().unwrap();
let v1: i64 = conn
.query_row(
"SELECT archived FROM versions WHERE package = ? AND version = ?",
params!["demo", "1.0.0"],
|r| r.get(0),
)
.unwrap();
let v2: i64 = conn
.query_row(
"SELECT archived FROM versions WHERE package = ? AND version = ?",
params!["demo", "2.0.0"],
|r| r.get(0),
)
.unwrap();
assert_eq!(v1, 1);
assert_eq!(v2, 0);
let _ = std::fs::remove_file(&src1);
let _ = std::fs::remove_file(&src2);
}
}
+170
View File
@@ -0,0 +1,170 @@
//! Black-box integration tests for the `ccc` binary's CLI subcommands
//! (create/add/archive). Each test runs the real compiled binary as a
//! subprocess against its own throwaway SQLite file, so no in-process
//! state is shared between tests.
use std::io::Write;
use std::path::PathBuf;
use std::process::Command;
fn bin() -> &'static str {
env!("CARGO_BIN_EXE_CCC")
}
struct TestEnv {
db_path: PathBuf,
}
impl TestEnv {
fn new(tag: &str) -> Self {
let db_path = std::env::temp_dir().join(format!(
"ccc-cli-test-{tag}-{}-{:?}.db",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let _ = std::fs::remove_file(&db_path);
Self { db_path }
}
fn cmd(&self, args: &[&str]) -> std::process::Output {
Command::new(bin())
.args(args)
.env("CCC_DB_PATH", &self.db_path)
.output()
.expect("failed to run ccc binary")
}
}
impl Drop for TestEnv {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.db_path);
}
}
fn write_temp_file(name: &str, contents: &[u8]) -> PathBuf {
// Keep `name` (with its extension) as the trailing path component so
// the CLI's mime-guessing logic sees the real file extension.
let dir = std::env::temp_dir().join(format!(
"ccc-cli-src-{}-{:?}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join(name);
let mut f = std::fs::File::create(&path).unwrap();
f.write_all(contents).unwrap();
path
}
#[test]
fn no_args_prints_usage_and_errors() {
let out = Command::new(bin())
.args(["bogus-subcommand"])
.output()
.expect("failed to run ccc binary");
assert!(!out.status.success());
let stderr = String::from_utf8_lossy(&out.stderr);
assert!(stderr.contains("unknown subcommand"), "stderr was: {stderr}");
}
#[test]
fn create_package_succeeds() {
let env = TestEnv::new("create-ok");
let out = env.cmd(&["create", "demo"]);
assert!(out.status.success(), "stderr: {}", String::from_utf8_lossy(&out.stderr));
let stdout = String::from_utf8_lossy(&out.stdout);
assert!(stdout.contains("created package 'demo'"), "stdout was: {stdout}");
}
#[test]
fn create_duplicate_package_fails() {
let env = TestEnv::new("create-dup");
let first = env.cmd(&["create", "demo"]);
assert!(first.status.success());
let second = env.cmd(&["create", "demo"]);
assert!(!second.status.success());
let stderr = String::from_utf8_lossy(&second.stderr);
assert!(stderr.contains("already exists"), "stderr was: {stderr}");
}
#[test]
fn add_to_missing_package_fails() {
let env = TestEnv::new("add-missing");
let src = write_temp_file("app.js", b"console.log(1)");
let out = env.cmd(&["add", "demo", src.to_str().unwrap(), "1.0.0"]);
assert!(!out.status.success());
let stderr = String::from_utf8_lossy(&out.stderr);
assert!(stderr.contains("does not exist"), "stderr was: {stderr}");
let _ = std::fs::remove_file(&src);
}
#[test]
fn add_missing_source_file_fails() {
let env = TestEnv::new("add-nofile");
assert!(env.cmd(&["create", "demo"]).status.success());
let out = env.cmd(&["add", "demo", "/nonexistent/path/does-not-exist.js", "1.0.0"]);
assert!(!out.status.success());
}
#[test]
fn full_lifecycle_create_add_archive() {
let env = TestEnv::new("lifecycle");
let src = write_temp_file("app.js", b"console.log('hello world');");
assert!(env.cmd(&["create", "demo"]).status.success());
let add_out = env.cmd(&["add", "demo", src.to_str().unwrap(), "1.0.0"]);
assert!(add_out.status.success(), "stderr: {}", String::from_utf8_lossy(&add_out.stderr));
let archive_version = env.cmd(&["archive", "demo", "1.0.0"]);
assert!(archive_version.status.success());
let archive_package = env.cmd(&["archive", "demo"]);
assert!(archive_package.status.success());
let archive_again = env.cmd(&["archive", "demo", "9.9.9"]);
assert!(!archive_again.status.success(), "archiving an unknown version should fail");
let _ = std::fs::remove_file(&src);
}
#[test]
fn stats_on_empty_db_reports_zeros() {
let env = TestEnv::new("stats-empty");
let out = env.cmd(&["stats"]);
assert!(out.status.success(), "stderr: {}", String::from_utf8_lossy(&out.stderr));
let stdout = String::from_utf8_lossy(&out.stdout);
assert!(stdout.contains("0 active, 0 archived"), "stdout was: {stdout}");
assert!(stdout.contains("0 bytes (gzipped)"), "stdout was: {stdout}");
}
#[test]
fn stats_reflects_added_content() {
let env = TestEnv::new("stats-content");
let src = write_temp_file("app.js", b"console.log('hello');");
assert!(env.cmd(&["create", "demo"]).status.success());
assert!(env.cmd(&["add", "demo", src.to_str().unwrap(), "1.0.0"]).status.success());
let out = env.cmd(&["stats"]);
assert!(out.status.success());
let stdout = String::from_utf8_lossy(&out.stdout);
assert!(stdout.contains("1 active, 0 archived"), "stdout was: {stdout}");
assert!(!stdout.contains("0 bytes (gzipped)"), "stdout was: {stdout}");
let _ = std::fs::remove_file(&src);
}
#[test]
fn missing_required_args_exit_with_usage_code() {
let env = TestEnv::new("usage");
let out = env.cmd(&["create"]);
assert_eq!(out.status.code(), Some(2));
let out = env.cmd(&["add", "demo"]);
assert_eq!(out.status.code(), Some(2));
}
+252
View File
@@ -0,0 +1,252 @@
//! End-to-end smoke test that actually boots `ccc serve` on the real
//! smarm/urus runtime and drives it over a raw TCP socket. This is the
//! regression test for the segfault-on-startup issue: if the runtime
//! still crashes on boot or on first request, this test hangs/fails
//! instead of the bug only showing up in production.
use std::io::{Read, Write};
use std::net::{TcpListener, TcpStream};
use std::path::PathBuf;
use std::process::{Child, Command, Stdio};
use std::time::{Duration, Instant};
fn bin() -> &'static str {
env!("CARGO_BIN_EXE_CCC")
}
fn free_port() -> u16 {
let listener = TcpListener::bind("127.0.0.1:0").expect("failed to bind ephemeral port");
listener.local_addr().unwrap().port()
}
struct Server {
child: Child,
port: u16,
db_path: PathBuf,
}
impl Server {
fn start(tag: &str) -> Self {
let db_path = std::env::temp_dir().join(format!(
"ccc-server-test-{tag}-{}-{:?}.db",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let _ = std::fs::remove_file(&db_path);
let port = free_port();
let child = Command::new(bin())
.args(["serve", "-p", &port.to_string()])
.env("CCC_DB_PATH", &db_path)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("failed to spawn ccc serve");
let server = Server { child, port, db_path };
server.wait_for_ready();
server
}
fn wait_for_ready(&self) {
let deadline = Instant::now() + Duration::from_secs(10);
loop {
if Instant::now() > deadline {
panic!("server on port {} did not become ready in time (possible segfault/hang on boot)", self.port);
}
match TcpStream::connect(("127.0.0.1", self.port)) {
Ok(_) => return,
Err(_) => std::thread::sleep(Duration::from_millis(50)),
}
}
}
fn cli(&self, args: &[&str]) -> std::process::Output {
Command::new(bin())
.args(args)
.env("CCC_DB_PATH", &self.db_path)
.output()
.expect("failed to run ccc CLI")
}
/// Sends a bare-bones HTTP/1.1 GET request and returns (status, headers-lowercased, body).
fn get(&self, path: &str) -> (u16, Vec<(String, String)>, Vec<u8>) {
let mut stream = TcpStream::connect(("127.0.0.1", self.port)).expect("connect failed");
stream.set_read_timeout(Some(Duration::from_secs(5))).unwrap();
let req = format!(
"GET {path} HTTP/1.1\r\nHost: 127.0.0.1\r\nAccept-Encoding: identity\r\nConnection: close\r\n\r\n"
);
stream.write_all(req.as_bytes()).expect("write failed");
let mut buf = Vec::new();
stream.read_to_end(&mut buf).expect("read failed");
parse_http_response(&buf)
}
/// Like `get`, but advertises gzip support so we can check the raw compressed path too.
fn get_gzip(&self, path: &str) -> (u16, Vec<(String, String)>, Vec<u8>) {
let mut stream = TcpStream::connect(("127.0.0.1", self.port)).expect("connect failed");
stream.set_read_timeout(Some(Duration::from_secs(5))).unwrap();
let req = format!(
"GET {path} HTTP/1.1\r\nHost: 127.0.0.1\r\nAccept-Encoding: gzip\r\nConnection: close\r\n\r\n"
);
stream.write_all(req.as_bytes()).expect("write failed");
let mut buf = Vec::new();
stream.read_to_end(&mut buf).expect("read failed");
parse_http_response(&buf)
}
}
impl Drop for Server {
fn drop(&mut self) {
let _ = self.child.kill();
let _ = self.child.wait();
let _ = std::fs::remove_file(&self.db_path);
}
}
fn parse_http_response(buf: &[u8]) -> (u16, Vec<(String, String)>, Vec<u8>) {
let sep = b"\r\n\r\n";
let split_at = buf
.windows(sep.len())
.position(|w| w == sep)
.expect("response missing header/body separator");
let head = std::str::from_utf8(&buf[..split_at]).expect("head not valid utf8");
let body = buf[split_at + sep.len()..].to_vec();
let mut lines = head.split("\r\n");
let status_line = lines.next().expect("missing status line");
let status: u16 = status_line
.split_whitespace()
.nth(1)
.expect("malformed status line")
.parse()
.expect("status code not numeric");
let headers = lines
.filter_map(|l| l.split_once(':'))
.map(|(k, v)| (k.trim().to_lowercase(), v.trim().to_string()))
.collect();
(status, headers, body)
}
fn header<'a>(headers: &'a [(String, String)], name: &str) -> Option<&'a str> {
headers.iter().find(|(k, _)| k == name).map(|(_, v)| v.as_str())
}
fn write_temp_file(name: &str, contents: &[u8]) -> PathBuf {
// Keep `name` (with its extension) as the trailing path component so
// the CLI's mime-guessing logic sees the real file extension.
let dir = std::env::temp_dir().join(format!(
"ccc-server-src-{}-{:?}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join(name);
let mut f = std::fs::File::create(&path).unwrap();
f.write_all(contents).unwrap();
path
}
#[test]
fn server_boots_without_segfaulting_and_answers_health_ish_route() {
let mut server = Server::start("boot");
// If we got past Server::start() the process accepted a TCP connection,
// i.e. it did not segfault/panic during startup.
let (status, _headers, body) = server.get("/packages");
assert_eq!(status, 200, "body: {}", String::from_utf8_lossy(&body));
let text = String::from_utf8_lossy(&body);
assert_eq!(text, "[]", "fresh DB should list no packages, got: {text}");
// The child must still be alive (no crash-on-request) after serving it.
match server.child.try_wait() {
Ok(None) => {}
Ok(Some(status)) => panic!("server process exited unexpectedly: {status}"),
Err(e) => panic!("failed to poll child status: {e}"),
}
}
#[test]
fn server_serves_added_asset_uncompressed_and_gzip() {
let server = Server::start("asset");
assert!(server.cli(&["create", "demo"]).status.success());
let src = write_temp_file("app.js", b"console.log('hello from ccc');");
let add = server.cli(&["add", "demo", src.to_str().unwrap(), "1.0.0"]);
assert!(add.status.success(), "stderr: {}", String::from_utf8_lossy(&add.stderr));
let (status, headers, body) = server.get("/assets/demo/1.0.0/app.js");
assert_eq!(status, 200, "body: {}", String::from_utf8_lossy(&body));
assert_eq!(header(&headers, "content-type"), Some("application/javascript"));
assert_eq!(body, b"console.log('hello from ccc');");
let (gz_status, gz_headers, gz_body) = server.get_gzip("/assets/demo/1.0.0/app.js");
assert_eq!(gz_status, 200);
assert_eq!(header(&gz_headers, "content-encoding"), Some("gzip"));
// Decompress and confirm round-trip integrity through the actual HTTP path.
use flate2::read::GzDecoder;
let mut decoder = GzDecoder::new(&gz_body[..]);
let mut decompressed = Vec::new();
decoder.read_to_end(&mut decompressed).unwrap();
assert_eq!(decompressed, b"console.log('hello from ccc');");
let _ = std::fs::remove_file(&src);
}
#[test]
fn server_returns_404_for_unknown_asset() {
let server = Server::start("404");
let (status, _headers, body) = server.get("/assets/ghost/9.9.9/missing.js");
assert_eq!(status, 404, "body: {}", String::from_utf8_lossy(&body));
}
#[test]
fn server_lists_packages_created_out_of_band_via_cli() {
let server = Server::start("listing");
assert!(server.cli(&["create", "alpha"]).status.success());
assert!(server.cli(&["create", "beta"]).status.success());
let src = write_temp_file("a.txt", b"x");
assert!(server.cli(&["add", "alpha", src.to_str().unwrap(), "1.0.0"]).status.success());
let (status, _headers, body) = server.get("/packages");
assert_eq!(status, 200);
let text = String::from_utf8_lossy(&body);
assert!(text.contains("alpha"), "body: {text}");
assert!(text.contains("beta"), "body: {text}");
assert!(text.contains("1.0.0"), "body: {text}");
let _ = std::fs::remove_file(&src);
}
#[test]
fn archived_package_is_hidden_from_listing() {
let server = Server::start("archived-listing");
assert!(server.cli(&["create", "demo"]).status.success());
assert!(server.cli(&["archive", "demo"]).status.success());
let (status, _headers, body) = server.get("/packages");
assert_eq!(status, 200);
let text = String::from_utf8_lossy(&body);
assert!(!text.contains("demo"), "archived package leaked into listing: {text}");
}
#[test]
fn server_survives_multiple_sequential_requests() {
// Repeated request/response cycles against the same long-lived process
// are exactly the pattern that would surface a use-after-free / double
// free in a bespoke runtime, hence several round trips here rather than
// just one.
let server = Server::start("multi");
for _ in 0..10 {
let (status, _headers, _body) = server.get("/packages");
assert_eq!(status, 200);
}
}