Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 1 | //! Reading a service's D1 database near the caller, with D1's Sessions API, |
| 2 | //! and saying how long an RPC took. | |
| 3 | //! | |
| 4 | //! The caller chooses, per request, with the `x-d1-bookmark` header: | |
| 5 | //! | |
| 6 | //! | Header | Session | Reads go to | | |
| 7 | //! | --- | --- | --- | | |
| 8 | //! | absent | none (the plain binding) | the primary, as always | | |
| 9 | //! | `first-primary` | started on the primary | the primary, then any copy at least as new | | |
| 10 | //! | `first-unconstrained` | started anywhere | the nearest copy | | |
| 11 | //! | a bookmark | started at that bookmark | any copy at least as new as the bookmark | | |
| 12 | //! | |
| 13 | //! Writes always go to the primary. A session is sequentially consistent: | |
| 14 | //! what it wrote, it reads back. Its latest bookmark comes back in the | |
| 15 | //! response's `x-d1-bookmark` header, so the caller can start the next | |
| 16 | //! request where this one left off. Without the header nothing changes, | |
| 17 | //! so service-to-service calls ([`crate::call`]), queues and crons read the | |
| 18 | //! primary as before. docs/PERFORMANCE.md explains who sends what. | |
| 19 | //! | |
| 20 | //! Every response that goes through [`Served::finish`] also carries | |
| 21 | //! `server-timing: svc;dur=<ms>;desc="<how D1 was read>"`, which the site | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 22 | //! adds up per service for its own `Server-Timing` header. A service that |
| 23 | //! times its round trips with a [`Timing`] and answers through | |
| 24 | //! [`Served::finish_timed`] adds `db;dur=<ms>;desc="<n> round trips, <m> | |
| 25 | //! statements"` and `rpc;dur=<ms>;desc="<n> calls"` beside it: summed time | |
| 26 | //! spent waiting on its database and on other services. | |
| 27 | ||
| 28 | use std::cell::Cell; | |
| 29 | use std::future::Future; | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 30 | |
| 31 | use worker::wasm_bindgen::JsCast; | |
| 32 | use worker::{D1Database, D1DatabaseSession, Env, Request, Response, Result}; | |
| 33 | ||
| 34 | /// The header that carries a session's constraint or bookmark, both ways. | |
| 35 | pub const BOOKMARK: &str = "x-d1-bookmark"; | |
| 36 | ||
| 37 | /// Longest bookmark accepted; D1's are about 70 characters. | |
| 38 | const MAX_BOOKMARK: usize = 256; | |
| 39 | ||
| 40 | /// What a request's `x-d1-bookmark` header asks for: `None` for no session | |
| 41 | /// (the primary, as without the header), otherwise what `withSession` is | |
| 42 | /// given. Anything that is not a constraint or a well-formed bookmark | |
| 43 | /// starts on the primary: never staler than asked. | |
| 44 | pub fn constraint(header: Option<&str>) -> Option<String> { | |
| 45 | let value = header?.trim(); | |
| 46 | if value.is_empty() { | |
| 47 | return None; | |
| 48 | } | |
| 49 | if value == "first-primary" || value == "first-unconstrained" { | |
| 50 | return Some(value.to_owned()); | |
| 51 | } | |
| 52 | let well_formed = value.len() <= MAX_BOOKMARK | |
| 53 | && value | |
| 54 | .chars() | |
| 55 | .all(|c| c.is_ascii_alphanumeric() || c == '-'); | |
| 56 | Some(if well_formed { value.to_owned() } else { "first-primary".to_owned() }) | |
| 57 | } | |
| 58 | ||
| 59 | /// How one RPC read its database, for its response. | |
| 60 | pub struct Served { | |
| 61 | session: Option<D1DatabaseSession>, | |
| 62 | started: u64, | |
| 63 | } | |
| 64 | ||
| 65 | /// The database `binding` for an RPC `request`: a D1 session when the | |
| 66 | /// caller asked for one, the plain binding otherwise. The session is | |
| 67 | /// handed back as a [`D1Database`] so a service's code is the same either | |
| 68 | /// way; it answers `prepare` and `batch` (all a request path uses), not | |
| 69 | /// `exec`, `dump` or `withSession`. | |
| 70 | pub fn open(env: &Env, binding: &str, request: &Request) -> Result<(D1Database, Served)> { | |
| 71 | let started = crate::now_ms(); | |
| 72 | let db = env.d1(binding)?; | |
| 73 | let asked = constraint(request.headers().get(BOOKMARK)?.as_deref()); | |
| 74 | let Some(asked) = asked else { | |
| 75 | return Ok((db, Served { session: None, started })); | |
| 76 | }; | |
| 77 | let session = db.with_session(Some(&asked))?; | |
| 78 | // The same JavaScript object, seen as a database: D1Database's methods | |
| 79 | // are structural, so `prepare` and `batch` call the session's own. | |
| 80 | let as_database = D1Database::unchecked_from_js(AsRef::<worker::wasm_bindgen::JsValue>::as_ref(&session).clone()); | |
| 81 | Ok((as_database, Served { session: Some(session), started })) | |
| 82 | } | |
| 83 | ||
| 84 | impl Served { | |
| 85 | /// For a request that touches no database: only its timing. | |
| 86 | pub fn timing_only() -> Self { | |
| 87 | Served { session: None, started: crate::now_ms() } | |
| 88 | } | |
| 89 | ||
| 90 | /// Adds the session's bookmark and the time taken to `response`. | |
| 91 | /// Errors pass through untouched. | |
| 92 | pub fn finish(&self, response: Result<Response>) -> Result<Response> { | |
| 93 | let mut response = response?; | |
| 94 | let took = crate::now_ms().saturating_sub(self.started); | |
| 95 | let how = if self.session.is_some() { "session" } else { "primary" }; | |
| 96 | let headers = response.headers_mut(); | |
| 97 | headers.append("server-timing", &format!("svc;dur={took};desc=\"{how}\""))?; | |
| 98 | if let Some(session) = &self.session | |
| 99 | && let Ok(Some(bookmark)) = session.get_bookmark() | |
| 100 | { | |
| 101 | headers.set(BOOKMARK, &bookmark)?; | |
| 102 | } | |
| 103 | Ok(response) | |
| 104 | } | |
| 105 | } | |
| 106 | ||
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 107 | /// Where one RPC's time went: round trips to its database and calls to |
| 108 | /// other services, each summed. One per request: a Worker serves several | |
| 109 | /// requests at once on one thread, so this is never global. | |
| 110 | #[derive(Default)] | |
| 111 | pub struct Timing { | |
| 112 | db_ms: Cell<u64>, | |
| 113 | db_trips: Cell<u32>, | |
| 114 | db_statements: Cell<u32>, | |
| 115 | rpc_ms: Cell<u64>, | |
| 116 | rpc_calls: Cell<u32>, | |
| 117 | } | |
| 118 | ||
| 119 | impl Timing { | |
| 120 | /// Times one round trip to the database carrying `statements` | |
| 121 | /// statements (a batch is one round trip). | |
| 122 | pub async fn db<T>(&self, statements: u32, work: impl Future<Output = T>) -> T { | |
| 123 | let started = crate::now_ms(); | |
| 124 | let answer = work.await; | |
| 125 | self.db_ms.set(self.db_ms.get() + crate::now_ms().saturating_sub(started)); | |
| 126 | self.db_trips.set(self.db_trips.get() + 1); | |
| 127 | self.db_statements.set(self.db_statements.get() + statements); | |
| 128 | answer | |
| 129 | } | |
| 130 | ||
| 131 | /// Times one call to another service. | |
| 132 | pub async fn rpc<T>(&self, work: impl Future<Output = T>) -> T { | |
| 133 | let started = crate::now_ms(); | |
| 134 | let answer = work.await; | |
| 135 | self.rpc_ms.set(self.rpc_ms.get() + crate::now_ms().saturating_sub(started)); | |
| 136 | self.rpc_calls.set(self.rpc_calls.get() + 1); | |
| 137 | answer | |
| 138 | } | |
| 139 | ||
| 140 | /// The `Server-Timing` entries, or `None` when nothing was timed. | |
| 141 | pub fn header(&self) -> Option<String> { | |
| 142 | let mut parts = Vec::new(); | |
| 143 | if self.db_trips.get() > 0 { | |
| 144 | parts.push(format!( | |
| 145 | "db;dur={};desc=\"{} round trips, {} statements\"", | |
| 146 | self.db_ms.get(), | |
| 147 | self.db_trips.get(), | |
| 148 | self.db_statements.get() | |
| 149 | )); | |
| 150 | } | |
| 151 | if self.rpc_calls.get() > 0 { | |
| 152 | parts.push(format!("rpc;dur={};desc=\"{} calls\"", self.rpc_ms.get(), self.rpc_calls.get())); | |
| 153 | } | |
| 154 | (!parts.is_empty()).then(|| parts.join(", ")) | |
| 155 | } | |
| 156 | } | |
| 157 | ||
| 158 | impl Served { | |
| 159 | /// [`Self::finish`], with what `timing` recorded. | |
| 160 | pub fn finish_timed(&self, response: Result<Response>, timing: &Timing) -> Result<Response> { | |
| 161 | let mut response = self.finish(response)?; | |
| 162 | if let Some(header) = timing.header() { | |
| 163 | response.headers_mut().append("server-timing", &header)?; | |
| 164 | } | |
| 165 | Ok(response) | |
| 166 | } | |
| 167 | } | |
| 168 | ||
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 169 | #[cfg(test)] |
| 170 | mod tests { | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 171 | use super::{Timing, constraint}; |
| 172 | ||
| 173 | #[test] | |
| 174 | fn nothing_timed_adds_nothing() { | |
| 175 | assert_eq!(Timing::default().header(), None); | |
| 176 | } | |
| 177 | ||
| 178 | #[test] | |
| 179 | fn round_trips_and_calls_are_summed() { | |
| 180 | let timing = Timing::default(); | |
| 181 | timing.db_ms.set(30); | |
| 182 | timing.db_trips.set(2); | |
| 183 | timing.db_statements.set(21); | |
| 184 | timing.rpc_ms.set(40); | |
| 185 | timing.rpc_calls.set(1); | |
| 186 | assert_eq!( | |
| 187 | timing.header().as_deref(), | |
| 188 | Some("db;dur=30;desc=\"2 round trips, 21 statements\", rpc;dur=40;desc=\"1 calls\"") | |
| 189 | ); | |
| 190 | } | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 191 | |
| 192 | #[test] | |
| 193 | fn no_header_means_the_primary_without_a_session() { | |
| 194 | assert_eq!(constraint(None), None); | |
| 195 | assert_eq!(constraint(Some("")), None); | |
| 196 | assert_eq!(constraint(Some(" ")), None); | |
| 197 | } | |
| 198 | ||
| 199 | #[test] | |
| 200 | fn constraints_and_bookmarks_pass_through() { | |
| 201 | assert_eq!(constraint(Some("first-primary")).as_deref(), Some("first-primary")); | |
| 202 | assert_eq!(constraint(Some("first-unconstrained")).as_deref(), Some("first-unconstrained")); | |
| 203 | let bookmark = "0000002c-00000004-00004f95-c7f4a9b2e8d1f0c3b6a5d4e3f2a1b0c9"; | |
| 204 | assert_eq!(constraint(Some(bookmark)).as_deref(), Some(bookmark)); | |
| 205 | } | |
| 206 | ||
| 207 | #[test] | |
| 208 | fn anything_else_starts_on_the_primary() { | |
| 209 | assert_eq!(constraint(Some("x; DROP")).as_deref(), Some("first-primary")); | |
| 210 | assert_eq!(constraint(Some(&"a".repeat(300))).as_deref(), Some("first-primary")); | |
| 211 | assert_eq!(constraint(Some("first primary")).as_deref(), Some("first-primary")); | |
| 212 | } | |
| 213 | } |