g1t/crates/kit/src/d1.rs

213 lines8,539 bytesCodeBlame

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 incidents1//! 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 daily22//! 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
28use std::cell::Cell;
29use std::future::Future;
Fast pages, required checks on the branch, self-hosted runners, honest incidents30
31use worker::wasm_bindgen::JsCast;
32use worker::{D1Database, D1DatabaseSession, Env, Request, Response, Result};
33
34/// The header that carries a session's constraint or bookmark, both ways.
35pub const BOOKMARK: &str = "x-d1-bookmark";
36
37/// Longest bookmark accepted; D1's are about 70 characters.
38const 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.
44pub 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.
60pub 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`.
70pub 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
84impl 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 daily107/// 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)]
111pub 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
119impl 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
158impl 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 incidents169#[cfg(test)]
170mod tests {
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily171 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 incidents191
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}