Initial commit: C3 (Cached Content Conduit)
A minimal CDN for serving versioned, gzip-compressed static assets, backed by SQLite and an actor-based (smarm) storage server, with HTTP serving via urus. - CLI: create/add/archive packages and versions, serve over HTTP - Storage: SQLite-backed asset store (src/store.rs) with gzip compression on ingest and content-type sniffing by extension - HTTP: GET /packages (list packages+versions), GET /assets/:package/:version/:filename (serves gzip or transparently decompressed, with immutable long-lived cache headers)
This commit is contained in:
+291
@@ -0,0 +1,291 @@
|
||||
mod store;
|
||||
|
||||
use rusqlite::{params, Connection};
|
||||
use std::io::Read;
|
||||
use std::sync::OnceLock;
|
||||
use flate2::read::GzDecoder;
|
||||
use serde::Serialize;
|
||||
use urus::{Config, Conn, Next, Pipeline, Router, serve_with};
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct AssetPayload {
|
||||
pub gzipped_bytes: Vec<u8>,
|
||||
pub mime_type: String,
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
pub struct PackageListing {
|
||||
name: String,
|
||||
versions: Vec<String>,
|
||||
}
|
||||
|
||||
pub enum GenServerMsg {
|
||||
FetchAsset {
|
||||
package: String,
|
||||
version: String,
|
||||
filename: String,
|
||||
reply_to: smarm::Sender<Result<Option<AssetPayload>, String>>,
|
||||
},
|
||||
ListPackages {
|
||||
reply_to: smarm::Sender<Vec<PackageListing>>,
|
||||
},
|
||||
}
|
||||
|
||||
pub struct AssetStoreServer {
|
||||
conn: Connection,
|
||||
}
|
||||
|
||||
impl AssetStoreServer {
|
||||
pub fn new() -> Self {
|
||||
let conn = store::open().expect("Failed to open SQLite database");
|
||||
Self { conn }
|
||||
}
|
||||
|
||||
pub fn loop_runner(self, rx: smarm::Receiver<GenServerMsg>) {
|
||||
while let Ok(msg) = rx.recv() {
|
||||
match msg {
|
||||
GenServerMsg::FetchAsset { package, version, filename, reply_to } => {
|
||||
let lookup = self.fetch_asset(&package, &version, &filename);
|
||||
let _ = reply_to.send(lookup);
|
||||
}
|
||||
GenServerMsg::ListPackages { reply_to } => {
|
||||
let listing = self.list_packages();
|
||||
let _ = reply_to.send(listing);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn fetch_asset(&self, package: &str, version: &str, filename: &str) -> Result<Option<AssetPayload>, String> {
|
||||
let mut stmt = self.conn
|
||||
.prepare(
|
||||
"SELECT gzipped_bytes, mime_type FROM versions
|
||||
WHERE package = ? AND version = ? AND filename = ?",
|
||||
)
|
||||
.map_err(|e| e.to_string())?;
|
||||
|
||||
let mut rows = stmt.query(params![package, version, filename]).map_err(|e| e.to_string())?;
|
||||
|
||||
if let Some(row) = rows.next().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())?;
|
||||
Ok(Some(AssetPayload { gzipped_bytes, mime_type }))
|
||||
} else {
|
||||
Ok(None)
|
||||
}
|
||||
}
|
||||
|
||||
fn list_packages(&self) -> Vec<PackageListing> {
|
||||
let mut listings = Vec::new();
|
||||
let mut stmt = match self.conn.prepare(
|
||||
"SELECT name FROM packages WHERE archived = 0 ORDER BY name",
|
||||
) {
|
||||
Ok(s) => s,
|
||||
Err(_) => return listings,
|
||||
};
|
||||
let names: Vec<String> = match stmt.query_map([], |r| r.get::<_, String>(0)) {
|
||||
Ok(rows) => rows.filter_map(Result::ok).collect(),
|
||||
Err(_) => return listings,
|
||||
};
|
||||
|
||||
for name in names {
|
||||
let mut vstmt = match self.conn.prepare(
|
||||
"SELECT version FROM versions WHERE package = ? AND archived = 0 ORDER BY version",
|
||||
) {
|
||||
Ok(s) => s,
|
||||
Err(_) => continue,
|
||||
};
|
||||
let versions: Vec<String> = match vstmt.query_map(params![name], |r| r.get::<_, String>(0)) {
|
||||
Ok(rows) => rows.filter_map(Result::ok).collect(),
|
||||
Err(_) => Vec::new(),
|
||||
};
|
||||
listings.push(PackageListing { name, versions });
|
||||
}
|
||||
listings
|
||||
}
|
||||
}
|
||||
|
||||
fn print_usage() {
|
||||
eprintln!(
|
||||
"usage:\n\
|
||||
\x20 ccc create <package>\n\
|
||||
\x20 ccc add <package> <filepath> <version>\n\
|
||||
\x20 ccc archive <package> [<version>]\n\
|
||||
\x20 ccc serve [-p|--port <port>]\n"
|
||||
);
|
||||
}
|
||||
|
||||
fn run_cli() -> Option<u16> {
|
||||
let args: Vec<String> = std::env::args().collect();
|
||||
match args.get(1).map(String::as_str) {
|
||||
Some("create") => {
|
||||
let Some(package) = args.get(2) else {
|
||||
print_usage();
|
||||
std::process::exit(2);
|
||||
};
|
||||
match store::create_package(package) {
|
||||
Ok(()) => println!("created package '{package}'"),
|
||||
Err(e) => {
|
||||
eprintln!("error: {e}");
|
||||
std::process::exit(1);
|
||||
}
|
||||
}
|
||||
None
|
||||
}
|
||||
Some("add") => {
|
||||
let (Some(package), Some(filepath), Some(version)) =
|
||||
(args.get(2), args.get(3), args.get(4))
|
||||
else {
|
||||
print_usage();
|
||||
std::process::exit(2);
|
||||
};
|
||||
match store::add_version(package, filepath, version) {
|
||||
Ok(()) => println!("added {package}@{version} from {filepath}"),
|
||||
Err(e) => {
|
||||
eprintln!("error: {e}");
|
||||
std::process::exit(1);
|
||||
}
|
||||
}
|
||||
None
|
||||
}
|
||||
Some("archive") => {
|
||||
let Some(package) = args.get(2) else {
|
||||
print_usage();
|
||||
std::process::exit(2);
|
||||
};
|
||||
let version = args.get(3).map(String::as_str);
|
||||
match store::archive(package, version) {
|
||||
Ok(()) => match version {
|
||||
Some(v) => println!("archived {package}@{v}"),
|
||||
None => println!("archived package '{package}'"),
|
||||
},
|
||||
Err(e) => {
|
||||
eprintln!("error: {e}");
|
||||
std::process::exit(1);
|
||||
}
|
||||
}
|
||||
None
|
||||
}
|
||||
Some("serve") | None => {
|
||||
let mut port: u16 = 8333;
|
||||
let mut i = 2;
|
||||
while i < args.len() {
|
||||
match args[i].as_str() {
|
||||
"-p" | "--port" => {
|
||||
let Some(v) = args.get(i + 1).and_then(|s| s.parse().ok()) else {
|
||||
eprintln!("error: --port requires a numeric argument");
|
||||
std::process::exit(2);
|
||||
};
|
||||
port = v;
|
||||
i += 2;
|
||||
}
|
||||
other => {
|
||||
eprintln!("unknown option '{other}'");
|
||||
print_usage();
|
||||
std::process::exit(2);
|
||||
}
|
||||
}
|
||||
}
|
||||
Some(port)
|
||||
}
|
||||
Some(other) => {
|
||||
eprintln!("unknown subcommand '{other}'");
|
||||
print_usage();
|
||||
std::process::exit(2);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The store actor can only be spawned from inside `smarm::Runtime::run()`
|
||||
// (i.e. from inside a connection actor, once `serve_with` has booted the
|
||||
// runtime). Lazily spawn it on first request and cache the handle.
|
||||
static STORE_TX: OnceLock<smarm::Sender<GenServerMsg>> = OnceLock::new();
|
||||
|
||||
fn store_handle() -> &'static smarm::Sender<GenServerMsg> {
|
||||
STORE_TX.get_or_init(|| {
|
||||
let (tx, rx) = smarm::channel::<GenServerMsg>();
|
||||
smarm::spawn(move || {
|
||||
let server = AssetStoreServer::new();
|
||||
server.loop_runner(rx);
|
||||
});
|
||||
tx
|
||||
})
|
||||
}
|
||||
|
||||
fn list_packages_handler(c: Conn, _n: Next) -> Conn {
|
||||
let (reply_tx, reply_rx) = smarm::channel();
|
||||
if store_handle().send(GenServerMsg::ListPackages { reply_to: reply_tx }).is_err() {
|
||||
return c.put_status(500).put_body("Internal DB Actor Failure");
|
||||
}
|
||||
match reply_rx.recv() {
|
||||
Ok(listing) => {
|
||||
let body = serde_json::to_vec(&listing).unwrap_or_default();
|
||||
c.put_status(200)
|
||||
.put_header("content-type", "application/json")
|
||||
.put_body(body)
|
||||
}
|
||||
Err(_) => c.put_status(500).put_body("Internal DB Actor Failure"),
|
||||
}
|
||||
}
|
||||
|
||||
fn fetch_asset_handler(c: Conn, _n: Next) -> Conn {
|
||||
let package = c.params.get("package").unwrap_or_default().to_string();
|
||||
let version = c.params.get("version").unwrap_or_default().to_string();
|
||||
let filename = c.params.get("filename").unwrap_or_default().to_string();
|
||||
|
||||
let (reply_tx, reply_rx) = smarm::channel();
|
||||
|
||||
if store_handle().send(GenServerMsg::FetchAsset {
|
||||
package,
|
||||
version,
|
||||
filename,
|
||||
reply_to: reply_tx,
|
||||
}).is_err() {
|
||||
return c.put_status(500).put_body("Internal DB Actor Failure");
|
||||
}
|
||||
|
||||
match reply_rx.recv() {
|
||||
Ok(Ok(Some(asset))) => {
|
||||
let accepts_gzip = c.headers
|
||||
.get("accept-encoding")
|
||||
.map(|v| v.contains("gzip"))
|
||||
.unwrap_or(false);
|
||||
|
||||
if accepts_gzip {
|
||||
c.put_status(200)
|
||||
.put_header("content-type", &asset.mime_type)
|
||||
.put_header("content-encoding", "gzip")
|
||||
.put_header("cache-control", "public, max-age=31536000, immutable")
|
||||
.put_header("access-control-allow-origin", "*")
|
||||
.put_body(asset.gzipped_bytes)
|
||||
} else {
|
||||
let mut decoder = GzDecoder::new(&asset.gzipped_bytes[..]);
|
||||
let mut raw_bytes = Vec::new();
|
||||
if decoder.read_to_end(&mut raw_bytes).is_ok() {
|
||||
c.put_status(200)
|
||||
.put_header("content-type", &asset.mime_type)
|
||||
.put_header("cache-control", "public, max-age=31536000, immutable")
|
||||
.put_header("access-control-allow-origin", "*")
|
||||
.put_body(raw_bytes)
|
||||
} else {
|
||||
c.put_status(500).put_body("Decompression Error")
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(Ok(None)) => c.put_status(404).put_body("Asset Not Found"),
|
||||
_ => c.put_status(500).put_body("Database Error Encountered"),
|
||||
}
|
||||
}
|
||||
|
||||
fn main() {
|
||||
let Some(port) = run_cli() else { return };
|
||||
|
||||
let router = Router::new()
|
||||
.get("/packages", list_packages_handler)
|
||||
.get("/assets/:package/:version/:filename", fetch_asset_handler);
|
||||
|
||||
println!("C3 on 0.0.0.0:{port}...");
|
||||
let cfg = Config::new(format!("0.0.0.0:{port}").parse().unwrap());
|
||||
let pipeline = Pipeline::new().plug(router);
|
||||
serve_with(cfg, pipeline).unwrap();
|
||||
}
|
||||
Reference in New Issue
Block a user