lib.rsannotatedlib.rssource269 lines · 13.5 KB · raw

The responses an Aldebaran site makes, over http and axum's Body.

The Worker family (workers-rs with its axum feature) and the native family (axum on tokio) both answer with http::Response<axum::body::Body>, so the adapter for the two is one crate. What needs the Cloudflare runtime itself (the console, the environment) is aldebaran-worker's.

Every function here makes the headers a response should have and none that it should not: a content type and a cache policy are always arguments, never defaults, and an HTML page can only be made with no-transform.

11use aldebaran_assets::{Asset, Bundle};
12use aldebaran_gate::{Asked, Gate};
13use aldebaran_headers::Policy;
14use axum::body::{Body, Bytes};
15use futures_util::{Stream, StreamExt};
16use http::{HeaderMap, HeaderName, HeaderValue, Method, Response, StatusCode, Uri, header};

A response with a content type and a cache policy, which are never left to a default.

19pub fn respond(content_type: &'static str, cache: &str, body: impl Into<Body>) -> Response<Body> {
20    Response::builder()
21        .header(header::CONTENT_TYPE, content_type)
22        .header(header::CACHE_CONTROL, cache)
23        .body(body.into())
24        .expect("a content type and a cache policy are valid header values")
25}

An HTML page. The cache policy is aldebaran_headers::html_cache_control, which always says no-transform, so Cloudflare does not add its analytics script.

29pub fn html(page: String, max_age_seconds: u32) -> Response<Body> {
30    respond("text/html; charset=utf-8", &aldebaran_headers::html_cache_control(max_age_seconds), page)
31}

The text of Datastar events (PatchElements::sse, PatchSignals::sse) as the response to a Datastar request. no-cache: an answer to a click is never reused.

35pub fn events(body: String) -> Response<Body> {
36    respond(aldebaran_datastar::EVENT_STREAM, "no-cache", body)
37}

A long-lived Datastar response: every item is the text of whole events (PatchElements::sse, PatchSignals::sse) and is sent as it arrives. no-cache, like events. Nothing is sent while the stream is quiet; a native server wraps it in [kept_alive] first.

42pub fn event_stream<S>(stream: S) -> Response<Body>
43where
44    S: Stream<Item = String> + Send + 'static,
45{
46    let chunks = stream.map(|text| Ok::<Bytes, std::convert::Infallible>(Bytes::from(text)));
47    respond(aldebaran_datastar::EVENT_STREAM, "no-cache", Body::from_stream(chunks))
48}

The line sent when a stream has been quiet: a comment, which a browser ignores and a proxy counts as traffic. It is what axum's Sse::keep_alive sends by default.

52#[cfg(feature = "keep-alive")]
53pub const KEEP_ALIVE: &str = ":\n\n";

How long a stream may be quiet before [KEEP_ALIVE] is sent, as axum's KeepAlive::default has it.

56#[cfg(feature = "keep-alive")]
57pub const KEEP_ALIVE_AFTER: std::time::Duration = std::time::Duration::from_secs(15);

The stream, with [KEEP_ALIVE] sent whenever quiet_for passes without an item. Any item restarts the wait, and the keep-alive ends with the stream.

61#[cfg(feature = "keep-alive")]
62pub fn kept_alive<S>(stream: S, quiet_for: std::time::Duration) -> impl Stream<Item = String> + Send + 'static
63where
64    S: Stream<Item = String> + Send + 'static,
65{
66    futures_util::stream::unfold(Box::pin(stream), move |mut stream| async move {
67        match tokio::time::timeout(quiet_for, stream.next()).await {
68            Ok(Some(text)) => Some((text, stream)),
69            Ok(None) => None,
70            Err(_) => Some((KEEP_ALIVE.to_owned(), stream)),
71        }
72    })
73}

Datastar's script, kept for a day (it is not named by its content).

76pub fn datastar_script() -> Response<Body> {
77    respond("text/javascript; charset=utf-8", "public, max-age=86400", aldebaran_datastar::SCRIPT)
78}

An embedded file, with the cache policy its query earns (Asset::cache_control).

81pub fn asset(asset: &Asset, query: Option<&str>) -> Response<Body> {
82    respond(asset.content_type, asset.cache_control(query), asset.body)
83}

The answer for a path no route took: the bundle's file if it has one, else not_found.

86pub fn bundle_or_not_found(bundle: &Bundle, uri: &Uri, nothing_here: &'static str) -> Response<Body> {
87    match bundle.find(uri.path()) {
88        Some(found) => asset(found, uri.query()),
89        None => not_found(nothing_here),
90    }
91}

404, with nosniff and no caching, saying one plain sentence.

94pub fn not_found(message: &'static str) -> Response<Body> {
95    let mut response = respond("text/plain; charset=utf-8", "no-store", message);
96    *response.status_mut() = StatusCode::NOT_FOUND;
97    response.headers_mut().insert(header::X_CONTENT_TYPE_OPTIONS, HeaderValue::from_static("nosniff"));
98    response
99}

422, for a value that is not one of a closed set. The message is the site's, never the value.

102pub fn refused(why: &'static str) -> Response<Body> {
103    let mut response = respond("text/plain; charset=utf-8", "no-store", why);
104    *response.status_mut() = StatusCode::UNPROCESSABLE_ENTITY;
105    response
106}

405 with the methods the site answers.

109pub fn method_not_allowed(allow: &'static str, why: String) -> Response<Body> {
110    let mut response = respond("text/plain; charset=utf-8", "no-store", why);
111    *response.status_mut() = StatusCode::METHOD_NOT_ALLOWED;
112    response.headers_mut().insert(header::ALLOW, HeaderValue::from_static(allow));
113    response
114}

Puts the policy's headers on a response, replacing any it already has. Every response leaves through this, so none leaves without them.

118pub fn secure(headers: &mut HeaderMap, policy: &Policy) {
119    for (name, value) in policy.headers() {
120        headers.insert(HeaderName::from_static(name), HeaderValue::from_str(&value).expect("a policy is plain ASCII"));
121    }
122}

A request as the gate reads it. Sec-Fetch-Mode other than navigate marks a request a page made for itself; an Upgrade: websocket marks a socket. A request with no Host is read by its URI.

126pub fn asked<'a>(method: &Method, headers: &'a HeaderMap, uri: &'a Uri) -> Asked<'a> {
127    let host = headers.get(header::HOST).and_then(|h| h.to_str().ok()).or_else(|| uri.host()).unwrap_or_default();
128    Asked {
129        host,
130        is_read: matches!(*method, Method::GET | Method::HEAD),
131        target: uri.path_and_query().map_or("/", |t| t.as_str()),
132        socket: headers.get(header::UPGRADE).is_some_and(|v| v.as_bytes().eq_ignore_ascii_case(b"websocket")),
133        own: headers.get("sec-fetch-mode").is_some_and(|mode| mode != "navigate"),
134    }
135}

What the gate decided, as a response; None means route the request.

138pub fn gated(gate: Gate) -> Option<Response<Body>> {
139    match gate {
140        Gate::Pass => None,
141        Gate::Moved { to, status } => {
142            let mut response = respond("text/plain; charset=utf-8", aldebaran_gate::MOVED_CACHE_CONTROL, format!("Moved to {to}\n"));
143            *response.status_mut() = StatusCode::from_u16(status.code()).expect("a redirect status");
144            response.headers_mut().insert(header::LOCATION, HeaderValue::from_str(&to).expect("the target of a valid request is a valid header"));
145            Some(response)
146        }
147        Gate::Refused(why) => Some(method_not_allowed("GET, HEAD", why)),
148    }
149}
151#[cfg(test)]
152mod tests {
153    use super::*;
154    use aldebaran_gate::{Canonical, Status, Unrepeatable};
155
156    fn parts(method: Method, host: &str, target: &str, extra: &[(&str, &str)]) -> (Method, HeaderMap, Uri) {
157        let mut headers = HeaderMap::new();
158        headers.insert(header::HOST, host.parse().unwrap());
159        for (name, value) in extra {
160            headers.insert(HeaderName::from_bytes(name.as_bytes()).unwrap(), value.parse().unwrap());
161        }
162        (method, headers, target.parse().unwrap())
163    }
164
165    #[test]
166    fn a_page_cannot_be_made_without_no_transform() {
167        let response = html("<p>".into(), 300);
168        assert_eq!(response.headers()[header::CACHE_CONTROL], "public, max-age=300, no-transform");
169        assert_eq!(response.headers()[header::CONTENT_TYPE], "text/html; charset=utf-8");
170    }
171
172    #[test]
173    fn events_are_a_stream_that_is_never_reused() {
174        let response = events("event: x\n\n".into());
175        assert_eq!(response.headers()[header::CONTENT_TYPE], "text/event-stream");
176        assert_eq!(response.headers()[header::CACHE_CONTROL], "no-cache");
177    }
178
179    async fn text(response: Response<Body>) -> String {
180        let bytes = axum::body::to_bytes(response.into_body(), usize::MAX).await.unwrap();
181        String::from_utf8(bytes.to_vec()).unwrap()
182    }
183
184    #[tokio::test]
185    async fn an_event_stream_sends_each_item_as_written() {
186        let items = futures_util::stream::iter(vec!["event: a\ndata: x\n\n".to_owned(), "event: b\ndata: y\n\n".to_owned()]);
187        let response = event_stream(items);
188        assert_eq!(response.headers()[header::CONTENT_TYPE], "text/event-stream");
189        assert_eq!(response.headers()[header::CACHE_CONTROL], "no-cache");
190        assert_eq!(text(response).await, "event: a\ndata: x\n\nevent: b\ndata: y\n\n");
191    }
192
193    #[cfg(feature = "keep-alive")]
194    #[tokio::test(start_paused = true)]
195    async fn a_quiet_stream_is_kept_alive_and_an_item_restarts_the_wait() {
196        use std::time::Duration;
197        let (sender, receiver) = tokio::sync::mpsc::unbounded_channel::<String>();
198        let items = futures_util::stream::unfold(receiver, |mut r| async move { r.recv().await.map(|t| (t, r)) });
199        let mut kept = Box::pin(kept_alive(items, KEEP_ALIVE_AFTER));
200        sender.send("event: a\n\n".into()).unwrap();
201        assert_eq!(kept.next().await.unwrap(), "event: a\n\n");
202        // Nothing for 15 seconds: a comment line.
203        assert_eq!(kept.next().await.unwrap(), ":\n\n");
204        assert_eq!(kept.next().await.unwrap(), ":\n\n");
205        // An item arriving 5 seconds into a wait is passed on, and the next keep-alive is a full 15 seconds later.
206        let started = tokio::time::Instant::now();
207        let waiting = tokio::spawn(async move { (kept.next().await, kept.next().await, tokio::time::Instant::now()) });
208        tokio::time::sleep(Duration::from_secs(5)).await;
209        sender.send("event: b\n\n".into()).unwrap();
210        let (first, second, at) = waiting.await.unwrap();
211        assert_eq!(first.unwrap(), "event: b\n\n");
212        assert_eq!(second.unwrap(), ":\n\n");
213        assert_eq!(at - started, Duration::from_secs(20));
214        // The wire text of the keep-alive is what axum's Sse sends.
215        assert_eq!(KEEP_ALIVE, ":\n\n");
216    }
217
218    #[cfg(feature = "keep-alive")]
219    #[tokio::test(start_paused = true)]
220    async fn the_keep_alive_ends_with_the_stream() {
221        let items = futures_util::stream::iter(vec!["event: a\n\n".to_owned()]);
222        let all: Vec<String> = kept_alive(items, KEEP_ALIVE_AFTER).collect().await;
223        assert_eq!(all, vec!["event: a\n\n".to_owned()]);
224    }
225
226    #[test]
227    fn a_request_is_read_for_the_gate() {
228        let (m, h, u) = parts(Method::GET, "old.example", "/a?b=1", &[("sec-fetch-mode", "cors")]);
229        let a = asked(&m, &h, &u);
230        assert_eq!((a.host, a.is_read, a.target, a.socket, a.own), ("old.example", true, "/a?b=1", false, true));
231        let (m, h, u) = parts(Method::POST, "x", "/", &[("upgrade", "WebSocket"), ("sec-fetch-mode", "navigate")]);
232        let a = asked(&m, &h, &u);
233        assert_eq!((a.is_read, a.socket, a.own), (false, true, false));
234    }
235
236    #[test]
237    fn the_gate_becomes_a_redirect_or_a_405_with_the_headers_every_response_has() {
238        let site = Canonical::every_other_leads_to("https://site.example").unwrap();
239        let (m, h, u) = parts(Method::GET, "old.example", "/?mascot=bunny", &[]);
240        let mut moved = gated(site.gate(&asked(&m, &h, &u), Status::PermanentRedirect, Unrepeatable::Refuse)).expect("sent on");
241        secure(moved.headers_mut(), &Policy::own_origin_only());
242        assert_eq!(moved.status(), StatusCode::PERMANENT_REDIRECT);
243        assert_eq!(moved.headers()[header::LOCATION], "https://site.example/?mascot=bunny");
244        assert_eq!(moved.headers()[header::CACHE_CONTROL], "public, max-age=86400");
245        assert_eq!(moved.headers()[header::X_CONTENT_TYPE_OPTIONS], "nosniff");
246
247        let (m, h, u) = parts(Method::POST, "old.example", "/mood", &[]);
248        let told = gated(site.gate(&asked(&m, &h, &u), Status::PermanentRedirect, Unrepeatable::Refuse)).expect("told");
249        assert_eq!(told.status(), StatusCode::METHOD_NOT_ALLOWED);
250        assert_eq!(told.headers()[header::ALLOW], "GET, HEAD");
251
252        let (m, h, u) = parts(Method::GET, "site.example", "/", &[]);
253        assert!(gated(site.gate(&asked(&m, &h, &u), Status::PermanentRedirect, Unrepeatable::Refuse)).is_none());
254    }
255
256    #[test]
257    fn a_bundle_answers_by_path_and_everything_else_is_not_found() {
258        static FILES: Bundle = Bundle::new(&[Asset { path: "/a.css", content_type: "text/css; charset=utf-8", body: b"a{}" }]);
259        let found = bundle_or_not_found(&FILES, &"/a.css".parse().unwrap(), "no");
260        assert_eq!(found.status(), StatusCode::OK);
261        assert_eq!(found.headers()[header::CACHE_CONTROL], "public, max-age=86400");
262        let own = format!("/a.css?v={}", FILES.find("/a.css").unwrap().version());
263        assert!(bundle_or_not_found(&FILES, &own.parse().unwrap(), "no").headers()[header::CACHE_CONTROL].to_str().unwrap().contains("immutable"));
264        let missing = bundle_or_not_found(&FILES, &"/nope".parse().unwrap(), "Nothing here.");
265        assert_eq!(missing.status(), StatusCode::NOT_FOUND);
266        assert_eq!(missing.headers()[header::X_CONTENT_TYPE_OPTIONS], "nosniff");
267        assert_eq!(refused("no").status(), StatusCode::UNPROCESSABLE_ENTITY);
268    }
269}