feat: serve rendered pages over http, spartan and nex

This commit is contained in:
randogoth 2026-10-04 19:55:21 +03:00
parent 46f12b10da
commit 2f38b8b929
30 changed files with 2896 additions and 90 deletions

View file

@ -11,13 +11,27 @@ publish = false
name = "itsybitsy"
path = "src/main.rs"
[features]
default = ["http", "spartan", "nex", "gemtext", "wap"]
http = []
spartan = []
nex = []
gemtext = ["dep:itsybitsy-gemtext"]
wap = ["dep:itsybitsy-wap"]
[dependencies]
anyhow = "1"
clap = { version = "4.6", features = ["derive"] }
env_logger = { version = "0.11", default-features = false, features = ["auto-color", "humantime"] }
itsybitsy-core = { path = "../core" }
itsybitsy-gemtext = { path = "../gemtext", optional = true }
itsybitsy-wap = { path = "../wap", optional = true }
log = "0.4"
[dev-dependencies]
itsybitsy-core = { path = "../core" }
tempfile = "3"
toml = "1.1"
[lints.rust]
unsafe_code = "forbid"

View file

@ -5,36 +5,45 @@
//! listener on port 0 in-process.
pub mod cli;
pub mod proto;
pub mod registry;
pub mod serve;
use std::path::Path;
use std::sync::Arc;
use anyhow::{Context, Result};
use itsybitsy_core::config::{Checked, ServerConfig};
use itsybitsy_core::render::Registry;
use itsybitsy_core::siteset::SiteSet;
use crate::cli::Cli;
pub fn run(cli: &Cli) -> Result<()> {
let formats = registry::format_ids();
let registry = Arc::new(registry::build());
let (config, checked, sites) = load(&cli.config, &registry.ids(), registry.clone())?;
if cli.check {
return check(&cli.config, &formats);
report(&config, &checked, &sites);
return Ok(());
}
let _ = load(&cli.config, &formats)?;
anyhow::bail!("serving is not implemented yet; run with --check to validate the configuration")
serve::run(&config, registry, Arc::new(sites))
}
/// Validate a configuration and report what it resolved to.
///
/// Takes the format set rather than reading [`registry`] directly, so a test can
/// supply the ids a later milestone's registry will provide.
/// check a configuration against ids this build may not provide.
pub fn check(config_path: &Path, formats: &[&str]) -> Result<()> {
let (config, checked, sites) = load(config_path, formats)?;
let (config, checked, sites) = load(config_path, formats, Arc::new(registry::build()))?;
report(&config, &checked, &sites);
Ok(())
}
fn load(config_path: &Path, formats: &[&str]) -> Result<(ServerConfig, Checked, SiteSet)> {
fn load(
config_path: &Path,
formats: &[&str],
registry: Arc<Registry>,
) -> Result<(ServerConfig, Checked, SiteSet)> {
// No context on load: its errors already name the file.
let config = ServerConfig::load(config_path)?;
let checked = config
@ -42,7 +51,8 @@ fn load(config_path: &Path, formats: &[&str]) -> Result<(ServerConfig, Checked,
.with_context(|| format!("validating {}", config_path.display()))?;
// Opening every root here means a broken site fails at startup rather than
// on the first request that happens to name it.
let sites = SiteSet::build(&config)
let rendered: Vec<String> = checked.formats.iter().cloned().collect();
let sites = SiteSet::build(&config, registry, rendered)
.with_context(|| format!("opening the sites in {}", config_path.display()))?;
Ok((config, checked, sites))
}

View file

@ -5,6 +5,10 @@ use clap::Parser;
use itsybitsy::cli::Cli;
fn main() -> ExitCode {
// `info` by default so a running server says what it is serving; `RUST_LOG`
// overrides it, as it does across the rest of the tree.
env_logger::Builder::from_env(env_logger::Env::default().default_filter_or("info")).init();
let cli = Cli::parse();
match itsybitsy::run(&cli) {
Ok(()) => ExitCode::SUCCESS,

224
bin/src/proto/http.rs Normal file
View file

@ -0,0 +1,224 @@
//! HTTP/1.1, hand-rolled and deliberately narrow: `GET` and `HEAD`, origin-form
//! targets, one response per connection.
//!
//! A reverse proxy in front sends origin-form, so absolute-form is refused rather
//! than half-supported. There is no keep-alive, no chunked encoding and no
//! pipelining; `Connection: close` is always sent so a client knows it.
use std::io::{BufReader, Write};
use std::net::TcpStream;
use anyhow::Result;
use itsybitsy_core::site::{Resolution, Resource};
use crate::proto::{for_log, negotiate, read_line_capped};
use crate::serve::Listener;
const MAX_REQUEST_LINE: usize = 8192;
const MAX_HEADER_LINE: usize = 8192;
const MAX_HEADERS: usize = 64;
const MAX_HEADER_BYTES: usize = 16 * 1024;
pub fn serve(listener: &Listener, mut stream: TcpStream) -> Result<()> {
let mut reader = BufReader::new(stream.try_clone()?);
let Some(line) = read_line_capped(&mut reader, MAX_REQUEST_LINE) else {
return bad_request(&mut stream);
};
let parts: Vec<&str> = line.split_whitespace().collect();
let [method, target, version] = parts.as_slice() else {
return bad_request(&mut stream);
};
if !version.starts_with("HTTP/1.") {
return bad_request(&mut stream);
}
// Origin-form only. smolweb's `urlsplit` would have mishandled the others.
if !target.starts_with('/') {
return bad_request(&mut stream);
}
let mut host = None;
let mut accept = None;
let mut total = 0usize;
for index in 0.. {
let Some(header) = read_line_capped(&mut reader, MAX_HEADER_LINE) else {
return bad_request(&mut stream);
};
if header.is_empty() {
break;
}
total += header.len();
if index >= MAX_HEADERS || total > MAX_HEADER_BYTES {
return bad_request(&mut stream);
}
let Some((name, value)) = header.split_once(':') else { continue };
let value = value.trim().to_string();
match name.trim().to_ascii_lowercase().as_str() {
"host" => host = Some(value),
"accept" => accept = Some(value),
_ => {}
}
}
let head_only = *method == "HEAD";
if *method != "GET" && !head_only {
return respond(
&mut stream,
405,
"Method Not Allowed",
"text/plain; charset=utf-8",
b"Method not allowed\n",
&[("Allow", "GET, HEAD")],
// Not HEAD: that case does not reach here.
false,
);
}
// A missing Host on HTTP/1.1 is a malformed request, not a missing resource:
// 404 would imply the server looked somewhere.
let Some(host) = host else { return bad_request(&mut stream) };
let Some(site) = listener.site_for(&host) else {
log::info!("{} http unknown host {}", listener.name, for_log(&host));
return respond(
&mut stream,
404,
"Not Found",
"text/plain; charset=utf-8",
b"Not found\n",
&[],
head_only,
);
};
let (path, query) = target.split_once('?').map_or((*target, ""), |(p, q)| (p, q));
log::info!("{} http {} {}", listener.name, for_log(&host), for_log(target));
// `?format=` overrides negotiation, for testing without the hardware.
let override_format = query
.split('&')
.find_map(|pair| pair.strip_prefix("format="))
.map(|value| value.to_ascii_lowercase());
let format = match &override_format {
Some(id) if listener.formats.iter().any(|f| f == id) => id.clone(),
Some(_) => {
return respond(
&mut stream,
400,
"Bad Request",
"text/plain; charset=utf-8",
b"Unknown format\n",
&[],
head_only,
);
}
None => {
negotiate::choose(&listener.registry, &listener.formats, accept.as_deref()).to_string()
}
};
match site.resolve(path) {
Ok(Resolution::Found(Resource::Document { page, .. })) => {
let Some(body) = page.body(&format) else {
return respond(
&mut stream,
500,
"Internal Server Error",
"text/plain; charset=utf-8",
b"",
&[],
head_only,
);
};
let cache = page.settings.cache_control.map(|age| format!("max-age={age}"));
let mut headers: Vec<(&str, &str)> = Vec::new();
// The response body depends on Accept, so a shared cache must not
// serve one client's format to another.
headers.push(("Vary", "Accept"));
if let Some(cache) = &cache {
headers.push(("Cache-Control", cache));
}
respond(&mut stream, 200, "OK", listener.media_type(&format), body, &headers, head_only)
}
Ok(Resolution::Found(Resource::Raw { path, media_type })) => {
let meta = std::fs::metadata(&path)?;
write_head(&mut stream, 200, "OK", media_type, meta.len(), &[])?;
if !head_only {
let mut file = std::fs::File::open(&path)?;
// Streamed, not buffered: smolweb reads the whole file into
// memory on every request.
std::io::copy(&mut file, &mut stream)?;
}
Ok(())
}
Ok(Resolution::Redirect(location)) => respond(
&mut stream,
301,
"Moved Permanently",
"text/plain; charset=utf-8",
b"",
&[("Location", location.as_str())],
head_only,
),
Ok(Resolution::NotFound) => respond(
&mut stream,
404,
"Not Found",
"text/plain; charset=utf-8",
b"Not found\n",
&[],
head_only,
),
Err(err) => {
log::warn!("{} http {}: {err}", listener.name, for_log(path));
respond(
&mut stream,
500,
"Internal Server Error",
"text/plain; charset=utf-8",
b"",
&[],
head_only,
)
}
}
}
fn bad_request(stream: &mut TcpStream) -> Result<()> {
respond(stream, 400, "Bad Request", "text/plain; charset=utf-8", b"Bad request\n", &[], false)
}
fn respond(
stream: &mut TcpStream,
code: u16,
reason: &str,
media_type: &str,
body: &[u8],
headers: &[(&str, &str)],
head_only: bool,
) -> Result<()> {
write_head(stream, code, reason, media_type, body.len() as u64, headers)?;
// HEAD sends the headers a GET would, including the length, and no body.
if !head_only {
stream.write_all(body)?;
}
Ok(())
}
fn write_head(
stream: &mut TcpStream,
code: u16,
reason: &str,
media_type: &str,
length: u64,
headers: &[(&str, &str)],
) -> Result<()> {
let mut head = format!(
"HTTP/1.1 {code} {reason}\r\nContent-Type: {media_type}\r\nContent-Length: {length}\r\n"
);
for (name, value) in headers {
head.push_str(&format!("{name}: {value}\r\n"));
}
head.push_str("Connection: close\r\n\r\n");
stream.write_all(head.as_bytes())?;
Ok(())
}

104
bin/src/proto/mod.rs Normal file
View file

@ -0,0 +1,104 @@
//! Shared plumbing for the protocol listeners.
pub mod http;
pub mod negotiate;
pub mod nex;
pub mod spartan;
use std::io::{self, BufRead, BufReader, Read};
use std::net::TcpStream;
use std::time::Duration;
/// Apply read and write timeouts to an accepted connection, so a client that
/// stops talking cannot hold a thread indefinitely.
pub fn set_timeouts(stream: &TcpStream, secs: u64) -> io::Result<()> {
let timeout = Some(Duration::from_secs(secs));
stream.set_read_timeout(timeout)?;
stream.set_write_timeout(timeout)
}
/// Read one CRLF- or LF-terminated line, refusing anything longer than `cap`.
///
/// Returns `None` when the line exceeds the cap or the connection ends before a
/// terminator, both of which the caller answers with its bad-request status.
pub fn read_line_capped<R: Read>(reader: &mut BufReader<R>, cap: usize) -> Option<String> {
let mut line = Vec::new();
// One past the cap, so a line of exactly `cap` bytes is still accepted.
let mut limited = reader.take(cap as u64 + 1);
limited.read_until(b'\n', &mut line).ok()?;
if line.len() > cap || !line.ends_with(b"\n") {
return None;
}
while line.last().is_some_and(|b| *b == b'\n' || *b == b'\r') {
line.pop();
}
String::from_utf8(line).ok()
}
/// Render untrusted text safe to put in a log line.
///
/// A request line can carry control bytes, and a newline in a log is how one
/// request's text becomes what looks like another's entry.
pub fn for_log(value: &str) -> String {
let mut out = String::with_capacity(value.len());
for ch in value.chars().take(256) {
if ch.is_control() {
out.push_str(&format!("\\x{:02x}", ch as u32 & 0xff));
} else {
out.push(ch);
}
}
out
}
#[cfg(test)]
mod tests {
use super::*;
fn read(input: &[u8], cap: usize) -> Option<String> {
read_line_capped(&mut BufReader::new(input), cap)
}
#[test]
fn reads_a_line_with_either_terminator() {
assert_eq!(read(b"GET /\r\n", 64).as_deref(), Some("GET /"));
assert_eq!(read(b"GET /\n", 64).as_deref(), Some("GET /"));
assert_eq!(read(b"\n", 64).as_deref(), Some(""));
}
#[test]
fn stops_at_the_first_line() {
assert_eq!(read(b"one\ntwo\n", 64).as_deref(), Some("one"));
}
#[test]
fn refuses_a_line_longer_than_the_cap() {
// The cap counts the terminator, so this is the longest accepted line.
assert_eq!(read(b"abcd\n", 5).as_deref(), Some("abcd"));
assert_eq!(read(b"abcde\n", 5), None);
}
#[test]
fn refuses_an_unterminated_line() {
// Otherwise a client could send a partial request and have it served.
assert_eq!(read(b"GET /", 64), None);
assert_eq!(read(b"", 64), None);
}
#[test]
fn refuses_a_line_that_is_not_utf8() {
assert_eq!(read(b"\xff\xfe\n", 64), None);
}
#[test]
fn escapes_control_bytes_for_logging() {
// A newline here is how one request forges another's log entry.
assert_eq!(for_log("GET /\r\nInjected: yes"), "GET /\\x0d\\x0aInjected: yes");
assert_eq!(for_log("/ordinary/path"), "/ordinary/path");
}
#[test]
fn truncates_a_very_long_value_for_logging() {
assert_eq!(for_log(&"a".repeat(1000)).len(), 256);
}
}

176
bin/src/proto/negotiate.rs Normal file
View file

@ -0,0 +1,176 @@
//! Choosing an output format from an HTTP `Accept` header.
//!
//! The one invariant: a wildcard never selects a format the client did not name
//! outright. A browser sends `Accept: text/html, application/xhtml+xml;q=0.9,
//! */*;q=0.8`, and that `*/*` must not be read as willingness to receive WML.
//! Only a literal media type counts as a match.
use itsybitsy_core::render::Registry;
/// Pick a format id from `accept`, out of the listener's `formats`.
///
/// `formats` is in preference order; the last entry is the fallback used when the
/// client named nothing this listener can serve, which is the ordinary case for a
/// browser.
pub fn choose<'a>(registry: &Registry, formats: &'a [String], accept: Option<&str>) -> &'a str {
let fallback = formats.last().map(String::as_str).unwrap_or("");
let Some(accept) = accept else { return fallback };
let ranges = parse(accept);
let mut best: Option<(&'a str, f32, usize)> = None;
for (index, id) in formats.iter().enumerate() {
let Some(renderer) = registry.get(id) else { continue };
let media_type = base_type(renderer.media_type());
// Literal matches only: a wildcard range is ignored entirely.
let Some(quality) = ranges
.iter()
.filter(|(range, _)| *range == media_type)
.map(|(_, q)| *q)
.fold(None, |best: Option<f32>, q| Some(best.map_or(q, |b: f32| b.max(q))))
else {
continue;
};
if quality <= 0.0 {
continue;
}
// A tie on quality is broken by the listener's own preference order.
let better = best.is_none_or(|(_, best_q, best_index)| {
quality > best_q || (quality == best_q && index < best_index)
});
if better {
best = Some((formats[index].as_str(), quality, index));
}
}
best.map(|(id, _, _)| id).unwrap_or(fallback)
}
/// `(media type, quality)` pairs, with parameters other than `q` discarded.
fn parse(accept: &str) -> Vec<(String, f32)> {
accept
.split(',')
.filter_map(|range| {
let mut parts = range.split(';');
let media_type = parts.next()?.trim().to_ascii_lowercase();
if media_type.is_empty() {
return None;
}
let quality = parts
.find_map(|param| {
let (name, value) = param.split_once('=')?;
(name.trim().eq_ignore_ascii_case("q")).then(|| value.trim())
})
.and_then(|value| value.parse::<f32>().ok())
.unwrap_or(1.0);
Some((media_type, quality))
})
.collect()
}
/// A renderer's media type without its parameters, for comparison against a range.
fn base_type(media_type: &str) -> String {
media_type.split(';').next().unwrap_or(media_type).trim().to_ascii_lowercase()
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use itsybitsy_core::Error;
use itsybitsy_core::ir::Doc;
use itsybitsy_core::render::{RenderCtx, Rendered, Renderer};
use super::*;
struct Stub(&'static str, &'static str);
impl Renderer for Stub {
fn id(&self) -> &'static str {
self.0
}
fn media_type(&self) -> &'static str {
self.1
}
fn render(&self, _doc: &Doc, _ctx: &RenderCtx<'_>) -> Result<Rendered, Error> {
Ok(Rendered::body(Vec::new()))
}
}
fn registry() -> Registry {
let mut registry = Registry::new();
registry.insert(Arc::new(Stub("wml", "text/vnd.wap.wml; charset=utf-8"))).unwrap();
registry
.insert(Arc::new(Stub("xhtmlmp", "application/vnd.wap.xhtml+xml; charset=utf-8")))
.unwrap();
registry.insert(Arc::new(Stub("html", "text/html; charset=utf-8"))).unwrap();
registry
}
fn formats() -> Vec<String> {
["wml", "xhtmlmp", "html"].iter().map(|s| s.to_string()).collect()
}
#[track_caller]
fn choose_with(accept: Option<&str>) -> String {
choose(&registry(), &formats(), accept).to_string()
}
#[test]
fn a_wildcard_never_selects_a_format_the_client_did_not_name() {
// The invariant the whole module exists for.
assert_eq!(choose_with(Some("*/*")), "html");
assert_eq!(choose_with(Some("text/*")), "html");
assert_eq!(choose_with(Some("text/html,application/xhtml+xml;q=0.9,*/*;q=0.8")), "html");
}
#[test]
fn no_header_at_all_falls_back() {
assert_eq!(choose_with(None), "html");
assert_eq!(choose_with(Some("")), "html");
}
#[test]
fn a_literal_media_type_is_honoured() {
assert_eq!(choose_with(Some("text/vnd.wap.wml")), "wml");
assert_eq!(choose_with(Some("application/vnd.wap.xhtml+xml")), "xhtmlmp");
assert_eq!(choose_with(Some("text/html")), "html");
}
#[test]
fn a_media_type_with_parameters_still_matches() {
assert_eq!(choose_with(Some("text/vnd.wap.wml; charset=utf-8")), "wml");
assert_eq!(choose_with(Some("TEXT/VND.WAP.WML")), "wml");
}
#[test]
fn quality_decides_between_two_named_formats() {
assert_eq!(choose_with(Some("text/vnd.wap.wml;q=0.5,text/html;q=0.9")), "html");
assert_eq!(choose_with(Some("text/vnd.wap.wml;q=0.9,text/html;q=0.5")), "wml");
}
#[test]
fn a_tie_breaks_by_the_listeners_preference_order() {
assert_eq!(choose_with(Some("text/html,text/vnd.wap.wml")), "wml");
let reversed: Vec<String> =
["html", "xhtmlmp", "wml"].iter().map(|s| s.to_string()).collect();
assert_eq!(choose(&registry(), &reversed, Some("text/html,text/vnd.wap.wml")), "html");
}
#[test]
fn a_zero_quality_refuses_that_format() {
assert_eq!(choose_with(Some("text/vnd.wap.wml;q=0")), "html");
}
#[test]
fn a_format_this_listener_does_not_serve_is_not_chosen() {
let only_html = vec!["html".to_string()];
assert_eq!(choose(&registry(), &only_html, Some("text/vnd.wap.wml")), "html");
}
#[test]
fn a_malformed_header_falls_back_rather_than_failing() {
assert_eq!(choose_with(Some(";;;")), "html");
assert_eq!(choose_with(Some("text/html;q=abc")), "html");
}
}

56
bin/src/proto/nex.rs Normal file
View file

@ -0,0 +1,56 @@
//! Nex: a bare path in, bytes out, no status line and no headers.
//!
//! There is nothing to signal an error with, so a failure is a plain-text body,
//! and nothing to redirect with either, which is why resolution follows the
//! canonical target server-side rather than bouncing the client.
use std::io::{BufReader, Write};
use std::net::TcpStream;
use anyhow::Result;
use itsybitsy_core::site::{Resolution, Resource};
use crate::proto::{for_log, read_line_capped};
use crate::serve::Listener;
/// Nex requests are a single path; the cap is smolweb's.
const MAX_REQUEST: usize = 2048;
pub fn serve(listener: &Listener, mut stream: TcpStream) -> Result<()> {
let selector = {
let mut reader = BufReader::new(stream.try_clone()?);
read_line_capped(&mut reader, MAX_REQUEST)
};
let Some(selector) = selector else {
log::warn!("{}: unreadable or oversized request", listener.name);
stream.write_all(b"Bad request\n")?;
return Ok(());
};
// Nex carries no hostname, so the listener names its site outright and
// configuration validation has already proved it exists.
let site = listener.only_site();
let path = if selector.is_empty() { "/" } else { &selector };
log::info!("{} nex {}", listener.name, for_log(path));
match site.resolve_flat(path) {
Ok(Resolution::Found(Resource::Document { url, page })) => {
match page.body(&listener.formats[0]) {
Some(body) => stream.write_all(body)?,
None => stream.write_all(b"Internal error\n")?,
}
let _ = url;
}
Ok(Resolution::Found(Resource::Raw { path, .. })) => {
let mut file = std::fs::File::open(&path)?;
std::io::copy(&mut file, &mut stream)?;
}
// resolve_flat never yields a redirect; it is matched for completeness.
Ok(_) => stream.write_all(b"Not found\n")?,
Err(err) => {
log::warn!("{} nex {}: {err}", listener.name, for_log(path));
stream.write_all(b"Internal error\n")?;
}
}
Ok(())
}

86
bin/src/proto/spartan.rs Normal file
View file

@ -0,0 +1,86 @@
//! Spartan: `HOST PATH LENGTH` in, a one-digit status and a body out.
//!
//! The host field is what smolweb discards; here it selects the virtual host,
//! falling back to the listener's `default_site` when it names nothing known.
use std::io::{BufReader, Read, Write};
use std::net::TcpStream;
use anyhow::Result;
use itsybitsy_core::site::{Resolution, Resource};
use crate::proto::{for_log, read_line_capped};
use crate::serve::Listener;
const MAX_REQUEST: usize = 4096;
/// An upload is refused either way, but a small body is drained first so the
/// connection stays in sync; a large one is not worth reading to discard.
const MAX_DRAIN: u64 = 1024 * 1024;
pub fn serve(listener: &Listener, mut stream: TcpStream) -> Result<()> {
let mut reader = BufReader::new(stream.try_clone()?);
let Some(line) = read_line_capped(&mut reader, MAX_REQUEST) else {
log::warn!("{}: unreadable or oversized request", listener.name);
return status(&mut stream, 4, "Bad request", b"");
};
let parts: Vec<&str> = line.split(' ').collect();
let [host, path, length] = parts.as_slice() else {
log::warn!("{} spartan: malformed request {}", listener.name, for_log(&line));
return status(&mut stream, 4, "Bad request", b"");
};
let Ok(length) = length.parse::<u64>() else {
return status(&mut stream, 4, "Bad request", b"");
};
if length > 0 {
// Drain before refusing, or the client's body would be read as the next
// request on a reused connection.
if length <= MAX_DRAIN {
let mut sink = io_sink();
std::io::copy(&mut reader.take(length), &mut sink)?;
}
return status(&mut stream, 4, "Uploads are not accepted", b"");
}
let Some(site) = listener.site_for(host) else {
log::info!("{} spartan unknown host {}", listener.name, for_log(host));
return status(&mut stream, 4, "Unknown host", b"");
};
log::info!("{} spartan {} {}", listener.name, for_log(host), for_log(path));
let format = &listener.formats[0];
match site.resolve(path) {
Ok(Resolution::Found(Resource::Document { page, .. })) => match page.body(format) {
Some(body) => {
let media_type = listener.media_type(format);
status(&mut stream, 2, media_type, body)
}
None => status(&mut stream, 5, "Not rendered", b""),
},
Ok(Resolution::Found(Resource::Raw { path, media_type })) => {
write!(stream, "2 {media_type}\r\n")?;
let mut file = std::fs::File::open(&path)?;
std::io::copy(&mut file, &mut stream)?;
Ok(())
}
Ok(Resolution::Redirect(location)) => status(&mut stream, 3, &location, b""),
Ok(Resolution::NotFound) => status(&mut stream, 4, "Not found", b""),
Err(err) => {
log::warn!("{} spartan {}: {err}", listener.name, for_log(path));
status(&mut stream, 5, "Internal error", b"")
}
}
}
fn status(stream: &mut TcpStream, code: u8, meta: &str, body: &[u8]) -> Result<()> {
write!(stream, "{code} {meta}\r\n")?;
if !body.is_empty() {
stream.write_all(body)?;
}
Ok(())
}
fn io_sink() -> std::io::Sink {
std::io::sink()
}

View file

@ -4,7 +4,21 @@
//! formats it was asked for. The registry is built once at startup and shared
//! immutably across every connection.
/// Ids of the formats this build provides.
pub fn format_ids() -> Vec<&'static str> {
Vec::new()
use std::sync::Arc;
use itsybitsy_core::render::Registry;
pub fn build() -> Registry {
let mut registry = Registry::new();
#[cfg(feature = "gemtext")]
registry.insert(Arc::new(itsybitsy_gemtext::Gemtext)).expect("duplicate renderer id");
#[cfg(feature = "wap")]
{
registry.insert(Arc::new(itsybitsy_wap::XhtmlMp)).expect("duplicate renderer id");
registry.insert(Arc::new(itsybitsy_wap::Html)).expect("duplicate renderer id");
}
registry
}

202
bin/src/serve.rs Normal file
View file

@ -0,0 +1,202 @@
//! Binding the listeners and dispatching connections.
//!
//! One thread per connection, with a per-listener cap. The work is read a short
//! line, stat a file, maybe render, write a buffer, close: blocking and
//! CPU-bound rather than connection-bound, which is the shape threads suit.
//! Async's advantage is many idle sockets, and this server has few busy ones.
use std::net::{TcpListener, TcpStream};
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
use std::thread;
use anyhow::{Context, Result};
use itsybitsy_core::config::{ListenerSpec, Protocol, ServerConfig};
use itsybitsy_core::render::Registry;
use itsybitsy_core::site::Site;
use itsybitsy_core::siteset::SiteSet;
use crate::proto;
/// Everything a connection handler needs, shared immutably across its threads.
pub struct Listener {
pub name: String,
pub protocol: Protocol,
/// Formats in preference order. Only HTTP has more than one.
pub formats: Vec<String>,
pub registry: Arc<Registry>,
pub sites: Arc<SiteSet>,
/// Served when a request names a host this server does not know.
pub default_site: Option<String>,
/// The single site served by a protocol that carries no hostname.
pub site: Option<String>,
pub timeout_secs: u64,
max_connections: u32,
open: AtomicU32,
}
impl Listener {
/// The site a request naming `host` is served from, falling back to
/// `default_site`.
pub fn site_for(&self, host: &str) -> Option<&Arc<Site>> {
self.sites
.lookup(host)
.or_else(|| self.default_site.as_deref().and_then(|name| self.sites.by_name(name)))
}
/// The one site a hostless protocol serves. Configuration validation proved
/// it exists, and with a single site defined the name may be left implicit.
pub fn only_site(&self) -> &Arc<Site> {
match &self.site {
Some(name) => self.sites.by_name(name).expect("validated at startup"),
None => self.sites.first().expect("at least one site is validated at startup"),
}
}
pub fn media_type(&self, format: &str) -> &'static str {
self.registry.get(format).map(|r| r.media_type()).unwrap_or("application/octet-stream")
}
}
/// A listener that has its socket but is not yet serving.
pub struct Bound {
pub listener: Arc<Listener>,
pub socket: TcpListener,
}
impl Bound {
pub fn local_addr(&self) -> std::io::Result<std::net::SocketAddr> {
self.socket.local_addr()
}
}
/// Bind every listener without serving any of them.
///
/// Separate from serving so that a port already in use is a startup failure
/// rather than one listener quietly missing, and so a test can learn the address
/// of a listener bound to port 0 before sending to it.
pub fn bind(
config: &ServerConfig,
registry: Arc<Registry>,
sites: Arc<SiteSet>,
) -> Result<Vec<Bound>> {
let mut bound = Vec::new();
for (name, spec) in &config.listener {
let socket = TcpListener::bind(&spec.bind)
.with_context(|| format!("listener '{name}' binding {}", spec.bind))?;
bound.push(Bound { listener: Arc::new(listener(name, spec, &registry, &sites)), socket });
}
Ok(bound)
}
/// Start one accept thread per bound listener.
pub fn spawn_all(bound: Vec<Bound>) -> Vec<thread::JoinHandle<()>> {
bound
.into_iter()
.map(|Bound { listener, socket }| thread::spawn(move || accept_loop(listener, socket)))
.collect()
}
/// Bind every listener, then serve until the process ends.
pub fn run(config: &ServerConfig, registry: Arc<Registry>, sites: Arc<SiteSet>) -> Result<()> {
let bound = bind(config, registry, sites)?;
for entry in &bound {
let local = entry.local_addr().map(|a| a.to_string()).unwrap_or_default();
log::info!(
"{} {} listening on {local}",
entry.listener.name,
entry.listener.protocol.as_str()
);
}
for thread in spawn_all(bound) {
// A panicking accept loop is a bug; the others keep serving.
if thread.join().is_err() {
log::error!("an accept loop panicked");
}
}
Ok(())
}
fn listener(
name: &str,
spec: &ListenerSpec,
registry: &Arc<Registry>,
sites: &Arc<SiteSet>,
) -> Listener {
Listener {
name: name.to_string(),
protocol: spec.protocol,
formats: spec.formats.clone(),
registry: registry.clone(),
sites: sites.clone(),
default_site: spec.default_site.clone(),
site: spec.site.clone(),
timeout_secs: spec.timeout_secs,
max_connections: spec.max_connections,
open: AtomicU32::new(0),
}
}
fn accept_loop(listener: Arc<Listener>, socket: TcpListener) {
for incoming in socket.incoming() {
let stream = match incoming {
Ok(stream) => stream,
Err(err) => {
log::warn!("{}: accept failed: {err}", listener.name);
continue;
}
};
// At the cap, refuse immediately rather than queueing threads without
// bound. smolweb has no cap at all.
let open = listener.open.fetch_add(1, Ordering::SeqCst);
if open >= listener.max_connections {
listener.open.fetch_sub(1, Ordering::SeqCst);
log::warn!("{}: at {} connections, refusing", listener.name, listener.max_connections);
continue;
}
let listener = listener.clone();
thread::spawn(move || {
handle(&listener, stream);
listener.open.fetch_sub(1, Ordering::SeqCst);
});
}
}
fn handle(listener: &Listener, stream: TcpStream) {
if let Err(err) = proto::set_timeouts(&stream, listener.timeout_secs) {
log::warn!("{}: setting timeouts failed: {err}", listener.name);
return;
}
// One malformed document must not take the process down, so a panic in a
// handler is caught and logged. This is the Rust equivalent of smolweb's
// bare `except Exception`, except the cause is recorded rather than lost.
let caught = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
match listener.protocol {
Protocol::Http => proto::http::serve(listener, stream),
Protocol::Spartan => proto::spartan::serve(listener, stream),
Protocol::Nex => proto::nex::serve(listener, stream),
// Validation refuses these until their listeners exist.
Protocol::Gemini | Protocol::Gopher => {
unreachable!("{} is rejected at startup", listener.protocol.as_str())
}
}
}));
match caught {
Ok(Ok(())) => {}
// A client that hung up mid-response is ordinary, not an error worth
// raising the level for.
Ok(Err(err)) => log::debug!("{}: {err}", listener.name),
Err(panic) => {
let reason = panic
.downcast_ref::<&str>()
.map(|s| s.to_string())
.or_else(|| panic.downcast_ref::<String>().cloned())
.unwrap_or_else(|| "unknown".to_string());
log::error!("{}: handler panicked: {reason}", listener.name);
}
}
}

403
bin/tests/listeners.rs Normal file
View file

@ -0,0 +1,403 @@
//! The three listeners over real sockets, bound on port 0 in-process.
//!
//! Two sites on one HTTP port, so virtual-host routing is exercised rather than
//! assumed.
use std::collections::BTreeMap;
use std::fs;
use std::io::{Read, Write};
use std::net::{SocketAddr, TcpStream};
use std::path::Path;
use std::sync::Arc;
use itsybitsy::serve;
use itsybitsy_core::config::ServerConfig;
use itsybitsy_core::siteset::SiteSet;
/// A running server, with the address of each listener by name.
struct Server {
addrs: BTreeMap<String, SocketAddr>,
_dir: tempfile::TempDir,
}
impl Server {
fn start() -> Self {
let dir = tempfile::tempdir().unwrap();
write(dir.path(), "one/index.md", "# One\n\nFirst site.\n");
write(dir.path(), "one/about.md", "# About\n\nAbout the first.\n");
write(dir.path(), "one/img.png", "not really a png");
write(dir.path(), "one/.itsybitsy.toml", "[defaults]\ncache_control = 120\n");
write(dir.path(), "two/index.md", "# Two\n\nSecond site.\n");
let text = format!(
r#"
version = 1
[site.one]
root = "{one}"
hosts = ["one.test"]
[site.two]
root = "{two}"
hosts = ["two.test"]
[listener.web]
protocol = "http"
bind = "127.0.0.1:0"
formats = ["xhtmlmp", "html"]
default_site = "one"
[listener.strict]
protocol = "http"
bind = "127.0.0.1:0"
formats = ["html"]
[listener.spartan]
protocol = "spartan"
bind = "127.0.0.1:0"
formats = ["gemtext"]
default_site = "one"
[listener.nex]
protocol = "nex"
bind = "127.0.0.1:0"
formats = ["gemtext"]
site = "two"
"#,
one = dir.path().join("one").to_str().unwrap(),
two = dir.path().join("two").to_str().unwrap(),
);
let config: ServerConfig = toml::from_str(&text).unwrap();
let registry = Arc::new(itsybitsy::registry::build());
let formats: Vec<String> =
["gemtext", "html", "xhtmlmp"].iter().map(|s| s.to_string()).collect();
let sites = Arc::new(SiteSet::build(&config, registry.clone(), formats).unwrap());
let bound = serve::bind(&config, registry, sites).unwrap();
let addrs = bound
.iter()
.map(|entry| (entry.listener.name.clone(), entry.local_addr().unwrap()))
.collect();
serve::spawn_all(bound);
Server { addrs, _dir: dir }
}
fn addr(&self, name: &str) -> SocketAddr {
self.addrs[name]
}
/// Send raw bytes to a listener and read the whole reply.
fn send(&self, listener: &str, request: &[u8]) -> String {
let mut stream = TcpStream::connect(self.addr(listener)).unwrap();
stream.write_all(request).unwrap();
stream.flush().unwrap();
let mut reply = Vec::new();
stream.read_to_end(&mut reply).unwrap();
String::from_utf8_lossy(&reply).into_owned()
}
fn get(&self, host: &str, target: &str) -> String {
self.request("web", host, target, "")
}
fn request(&self, listener: &str, host: &str, target: &str, extra: &str) -> String {
let request = format!("GET {target} HTTP/1.1\r\nHost: {host}\r\n{extra}\r\n");
self.send(listener, request.as_bytes())
}
}
fn write(root: &Path, rel: &str, body: &str) {
let path = root.join(rel);
fs::create_dir_all(path.parent().unwrap()).unwrap();
fs::write(path, body).unwrap();
}
/// Split a response into its status line and body.
fn split(response: &str) -> (&str, &str) {
let status = response.lines().next().unwrap_or("");
let body = response.split_once("\r\n\r\n").map(|(_, b)| b).unwrap_or("");
(status, body)
}
// -- HTTP ----------------------------------------------------------------
#[test]
fn serves_a_document() {
let server = Server::start();
let response = server.get("one.test", "/");
let (status, body) = split(&response);
assert_eq!(status, "HTTP/1.1 200 OK");
assert!(body.contains("<h1>One</h1>"), "{body}");
assert!(response.contains("Vary: Accept"), "{response}");
assert!(response.contains("Cache-Control: max-age=120"), "{response}");
}
#[test]
fn routes_two_sites_on_one_port() {
let server = Server::start();
assert!(server.get("one.test", "/").contains("<h1>One</h1>"));
assert!(server.get("two.test", "/").contains("<h1>Two</h1>"));
}
#[test]
fn an_unknown_host_falls_back_to_the_default_site() {
let server = Server::start();
assert!(server.get("elsewhere.test", "/").contains("<h1>One</h1>"));
}
#[test]
fn an_unknown_host_with_no_default_is_not_found() {
let server = Server::start();
let (status, body) = {
let response = server.request("strict", "elsewhere.test", "/", "");
(response.lines().next().unwrap().to_string(), response)
};
assert_eq!(status, "HTTP/1.1 404 Not Found");
assert!(body.ends_with("Not found\n"), "{body}");
}
#[test]
fn a_missing_host_header_is_a_bad_request_not_a_missing_page() {
// 404 would imply the server looked somewhere.
let server = Server::start();
let response = server.send("web", b"GET / HTTP/1.1\r\n\r\n");
assert_eq!(split(&response).0, "HTTP/1.1 400 Bad Request");
}
#[test]
fn the_md_form_redirects_to_the_canonical_url() {
let server = Server::start();
let response = server.get("one.test", "/about.md");
assert_eq!(split(&response).0, "HTTP/1.1 301 Moved Permanently");
assert!(response.contains("Location: /about"), "{response}");
}
#[test]
fn a_missing_page_is_not_found_with_a_fixed_body() {
// The body never echoes the request path.
let server = Server::start();
let response = server.get("one.test", "/nope");
let (status, body) = split(&response);
assert_eq!(status, "HTTP/1.1 404 Not Found");
assert_eq!(body, "Not found\n");
}
#[test]
fn a_raw_file_is_served_with_its_media_type_and_no_vary() {
let server = Server::start();
let response = server.get("one.test", "/img.png");
assert!(response.contains("Content-Type: image/png"), "{response}");
// Its bytes do not depend on Accept, so a shared cache need not vary on it.
assert!(!response.contains("Vary:"), "{response}");
}
#[test]
fn head_sends_the_headers_of_a_get_and_no_body() {
let server = Server::start();
let get = server.get("one.test", "/");
let head = server.send("web", b"HEAD / HTTP/1.1\r\nHost: one.test\r\n\r\n");
let length = |response: &str| {
response.lines().find_map(|line| line.strip_prefix("Content-Length: ")).unwrap().to_string()
};
assert_eq!(length(&get), length(&head));
assert_eq!(split(&head).1, "");
}
#[test]
fn an_unsupported_method_is_refused_with_the_allowed_set() {
let server = Server::start();
let response = server.send("web", b"POST / HTTP/1.1\r\nHost: one.test\r\n\r\n");
let (status, body) = split(&response);
assert_eq!(status, "HTTP/1.1 405 Method Not Allowed");
assert!(response.contains("Allow: GET, HEAD"), "{response}");
// The declared length has to match what is sent, or the response truncates.
assert_eq!(body, "Method not allowed\n");
}
#[test]
fn every_response_body_matches_its_declared_length() {
let server = Server::start();
let cases: Vec<Vec<u8>> = vec![
b"GET / HTTP/1.1\r\nHost: one.test\r\n\r\n".to_vec(),
b"GET /nope HTTP/1.1\r\nHost: one.test\r\n\r\n".to_vec(),
b"POST / HTTP/1.1\r\nHost: one.test\r\n\r\n".to_vec(),
b"GARBAGE\r\n\r\n".to_vec(),
b"GET /?format=bogus HTTP/1.1\r\nHost: one.test\r\n\r\n".to_vec(),
];
for request in cases {
let response = server.send("web", &request);
let (_, body) = split(&response);
let declared: usize = response
.lines()
.find_map(|line| line.strip_prefix("Content-Length: "))
.unwrap()
.parse()
.unwrap();
assert_eq!(
body.len(),
declared,
"for {:?} -> {response:?}",
String::from_utf8_lossy(&request)
);
}
}
#[test]
fn an_absolute_form_target_is_refused() {
// A reverse proxy in front sends origin-form.
let server = Server::start();
let response = server.send("web", b"GET http://one.test/ HTTP/1.1\r\nHost: one.test\r\n\r\n");
assert_eq!(split(&response).0, "HTTP/1.1 400 Bad Request");
}
#[test]
fn an_oversized_request_line_is_refused() {
let server = Server::start();
let mut request = b"GET /".to_vec();
request.extend(std::iter::repeat_n(b'a', 9000));
request.extend_from_slice(b" HTTP/1.1\r\nHost: one.test\r\n\r\n");
assert_eq!(split(&server.send("web", &request)).0, "HTTP/1.1 400 Bad Request");
}
#[test]
fn too_many_headers_are_refused() {
let server = Server::start();
let mut request = b"GET / HTTP/1.1\r\nHost: one.test\r\n".to_vec();
for index in 0..200 {
request.extend_from_slice(format!("X-Pad-{index}: x\r\n").as_bytes());
}
request.extend_from_slice(b"\r\n");
assert_eq!(split(&server.send("web", &request)).0, "HTTP/1.1 400 Bad Request");
}
// -- Negotiation ---------------------------------------------------------
#[test]
fn a_browser_accept_header_gets_html() {
let server = Server::start();
let response = server.request(
"web",
"one.test",
"/",
"Accept: text/html,application/xhtml+xml;q=0.9,*/*;q=0.8\r\n",
);
assert!(response.contains("Content-Type: text/html"), "{response}");
assert!(!split(&response).1.starts_with("<?xml"), "the html form carries no declaration");
}
#[test]
fn a_literal_wap_accept_header_gets_xhtml_mobile_profile() {
let server = Server::start();
let response =
server.request("web", "one.test", "/", "Accept: application/vnd.wap.xhtml+xml\r\n");
assert!(response.contains("Content-Type: application/vnd.wap.xhtml+xml"), "{response}");
assert!(split(&response).1.starts_with("<?xml version=\"1.0\""), "{response}");
}
#[test]
fn the_format_query_overrides_negotiation() {
let server = Server::start();
let response = server.get("one.test", "/?format=xhtmlmp");
assert!(response.contains("Content-Type: application/vnd.wap.xhtml+xml"), "{response}");
}
#[test]
fn an_unknown_format_query_is_a_bad_request() {
let server = Server::start();
let response = server.get("one.test", "/?format=bogus");
assert_eq!(split(&response).0, "HTTP/1.1 400 Bad Request");
}
// -- Spartan -------------------------------------------------------------
#[test]
fn spartan_serves_gemtext() {
let server = Server::start();
let response = server.send("spartan", b"one.test / 0\r\n");
assert!(response.starts_with("2 text/gemini; charset=utf-8\r\n"), "{response}");
assert!(response.contains("# One"), "{response}");
}
#[test]
fn spartan_routes_by_the_host_field_smolweb_discards() {
let server = Server::start();
assert!(server.send("spartan", b"one.test / 0\r\n").contains("# One"));
assert!(server.send("spartan", b"two.test / 0\r\n").contains("# Two"));
}
#[test]
fn spartan_redirects_the_md_form() {
let server = Server::start();
assert_eq!(server.send("spartan", b"one.test /about.md 0\r\n"), "3 /about\r\n");
}
#[test]
fn spartan_reports_a_missing_page() {
let server = Server::start();
assert_eq!(server.send("spartan", b"one.test /nope 0\r\n"), "4 Not found\r\n");
}
#[test]
fn spartan_drains_an_upload_then_refuses_it() {
// Draining keeps the connection in sync rather than leaving a body to be
// read as the next request.
let server = Server::start();
let response = server.send("spartan", b"one.test / 5\r\nhello");
assert_eq!(response, "4 Uploads are not accepted\r\n");
}
#[test]
fn spartan_refuses_a_malformed_request_line() {
let server = Server::start();
assert_eq!(server.send("spartan", b"only-two-fields 0\r\n"), "4 Bad request\r\n");
assert_eq!(server.send("spartan", b"one.test / notanumber\r\n"), "4 Bad request\r\n");
}
// -- Nex -----------------------------------------------------------------
#[test]
fn nex_sends_bytes_with_no_status_or_headers() {
let server = Server::start();
// This listener names site two, which has no hosts in its request at all.
let response = server.send("nex", b"/\r\n");
assert!(response.starts_with("# Two"), "{response}");
}
#[test]
fn nex_resolves_a_redirect_server_side() {
// There is no redirect status to bounce with, so the target's content comes
// back on the first request.
let server = Server::start();
let response = server.send("nex", b"/index.md\r\n");
assert!(response.starts_with("# Two"), "{response}");
}
#[test]
fn nex_reports_a_missing_page_as_a_plain_body() {
let server = Server::start();
assert_eq!(server.send("nex", b"/nope\r\n"), "Not found\n");
}
#[test]
fn an_empty_nex_selector_is_the_root() {
let server = Server::start();
assert!(server.send("nex", b"\r\n").starts_with("# Two"));
}
// -- Containment ---------------------------------------------------------
#[test]
fn nothing_outside_the_root_is_reachable_over_any_protocol() {
let server = Server::start();
for target in ["/../two/index.md", "/%2e%2e/two/index.md", "/.itsybitsy.toml"] {
let response = server.get("one.test", target);
let status = split(&response).0;
assert!(
status.starts_with("HTTP/1.1 404") || status.starts_with("HTTP/1.1 301"),
"{target} -> {status}"
);
assert!(!response.contains("<h1>Two</h1>"), "{target} reached the other site");
}
assert_eq!(server.send("spartan", b"one.test /.itsybitsy.toml 0\r\n"), "4 Not found\r\n");
assert_eq!(server.send("nex", b"/.itsybitsy.toml\r\n"), "Not found\n");
}