Skip to content
File

Blob: src/rust/kj/http.rs

rust526 lines
1use std::marker::PhantomData;
2use std::pin::Pin;
3 
4use futures::TryFutureExt;
5use kj_rs::KjOwn;
6use static_assertions::assert_eq_align;
7use static_assertions::assert_eq_size;
8 
9use crate::OwnOrMut;
10use crate::Result;
11use crate::io::AsyncInputStream;
12use crate::io::AsyncIoStream;
13 
14#[cxx::bridge(namespace = "kj::rust")]
15#[expect(clippy::missing_panics_doc)]
16#[expect(clippy::missing_safety_doc)]
17pub mod ffi {
18 unsafe extern "C++" {
19 include!("workerd/rust/kj/ffi.h");
20 }
21 
22 /// Corresponds to `kj::HttpMethod`.
23 /// Values are automatically assigned by `cxx` because of extern declaration below.
24 #[derive(Debug, PartialEq, Eq, Copy, Clone)]
25 #[repr(u32)]
26 enum HttpMethod {
27 GET,
28 HEAD,
29 POST,
30 PUT,
31 DELETE,
32 PATCH,
33 PURGE,
34 OPTIONS,
35 TRACE,
36 COPY,
37 LOCK,
38 MKCOL,
39 MOVE,
40 PROPFIND,
41 PROPPATCH,
42 SEARCH,
43 UNLOCK,
44 ACL,
45 REPORT,
46 MKACTIVITY,
47 CHECKOUT,
48 MERGE,
49 MSEARCH,
50 NOTIFY,
51 SUBSCRIBE,
52 UNSUBSCRIBE,
53 QUERY,
54 BAN,
55 }
56 unsafe extern "C++" {
57 type HttpMethod;
58 }
59 
60 // --- HttpHeaderId
61 // Opaque handle to a kj::HttpHeaderId, which identifies a header by numeric index in an
62 // HttpHeaderTable. This supports both builtin headers and custom headers registered via
63 // HttpHeaderTable::Builder::add(). Pass these by reference from C++ to Rust and back.
64 
65 unsafe extern "C++" {
66 type HttpHeaderId;
67 }
68 
69 // --- HttpHeaders
70 // TODO(when needed): support HttpHeaderId creation from rust.
71 
72 /// Corresponds to `kj::HttpHeaders::BuiltinIndicesEnum`.
73 /// Values are automatically assigned by `cxx` because of extern declaration below.
74 #[derive(Debug, PartialEq, Eq, Copy, Clone)]
75 #[repr(u32)]
76 pub enum BuiltinIndicesEnum {
77 CONNECTION,
78 KEEP_ALIVE,
79 TE,
80 TRAILER,
81 UPGRADE,
82 CONTENT_LENGTH,
83 TRANSFER_ENCODING,
84 SEC_WEBSOCKET_KEY,
85 SEC_WEBSOCKET_VERSION,
86 SEC_WEBSOCKET_ACCEPT,
87 SEC_WEBSOCKET_EXTENSIONS,
88 HOST,
89 DATE,
90 LOCATION,
91 CONTENT_TYPE,
92 RANGE,
93 CONTENT_RANGE,
94 }
95 
96 unsafe extern "C++" {
97 type BuiltinIndicesEnum;
98 type HttpHeaderTable;
99 type HttpHeaders;
100 fn new_http_headers(table: &HttpHeaderTable) -> KjOwn<HttpHeaders>;
101 fn clone_shallow(this_: &HttpHeaders) -> KjOwn<HttpHeaders>;
102 fn set_header(this_: Pin<&mut HttpHeaders>, id: BuiltinIndicesEnum, value: &str);
103 unsafe fn get_header<'a>(
104 this_: &'a HttpHeaders,
105 id: BuiltinIndicesEnum,
106 ) -> KjMaybe<&'a [u8]>;
107 unsafe fn get_header_by_id<'a>(
108 this_: &'a HttpHeaders,
109 id: &HttpHeaderId,
110 ) -> KjMaybe<&'a [u8]>;
111 }
112 
113 // --- kj::HttpService ffi
114 
115 unsafe extern "C++" {
116 type TlsStarterCallback;
117 }
118 
119 /// Corresponds to `kj::HttpConnectSettings`.
120 struct HttpConnectSettings<'a> {
121 use_tls: bool,
122 tls_starter: KjMaybe<Pin<&'a mut TlsStarterCallback>>,
123 }
124 
125 unsafe extern "C++" {
126 type AsyncInputStream = crate::io::ffi::AsyncInputStream;
127 type AsyncIoStream = crate::io::ffi::AsyncIoStream;
128 type AsyncOutputStream = crate::io::ffi::AsyncOutputStream;
129 type ConnectResponse;
130 type HttpServiceResponse;
131 type HttpService;
132 
133 fn response_send(
134 this_: Pin<&mut HttpServiceResponse>,
135 status_code: u32,
136 status_text: &str,
137 headers: &HttpHeaders,
138 expected_body_size: KjMaybe<u64>,
139 ) -> Result<KjOwn<AsyncOutputStream>>;
140 
141 fn connect_response_accept(
142 this_: Pin<&mut ConnectResponse>,
143 status_code: u32,
144 status_text: &str,
145 headers: &HttpHeaders,
146 ) -> Result<()>;
147 
148 fn connect_response_reject(
149 this_: Pin<&mut ConnectResponse>,
150 status_code: u32,
151 status_text: &str,
152 headers: &HttpHeaders,
153 expected_body_size: KjMaybe<u64>,
154 ) -> Result<KjOwn<AsyncOutputStream>>;
155 
156 /// Corresponds to `kj::HttpService::request`.
157 async fn request(
158 this_: Pin<&mut HttpService>,
159 method: HttpMethod,
160 url: &[u8],
161 headers: &HttpHeaders,
162 request_body: Pin<&mut AsyncInputStream>,
163 response: Pin<&mut HttpServiceResponse>,
164 ) -> Result<()>;
165 
166 /// Corresponds to `kj::HttpService::connect`.
167 async fn connect(
168 this_: Pin<&mut HttpService>,
169 host: &[u8],
170 headers: &HttpHeaders,
171 connection: Pin<&mut AsyncIoStream>,
172 response: Pin<&mut ConnectResponse>,
173 settings: HttpConnectSettings<'_>,
174 ) -> Result<()>;
175 }
176 
177 // DynHttpService
178 extern "Rust" {
179 type DynHttpService;
180 
181 async unsafe fn request<'a>(
182 self: &'a mut DynHttpService,
183 method: HttpMethod,
184 url: &'a [u8],
185 headers: &'a HttpHeaders,
186 request_body: Pin<&'a mut AsyncInputStream>,
187 response: Pin<&'a mut HttpServiceResponse>,
188 ) -> Result<()>;
189 
190 async unsafe fn connect<'a>(
191 self: &'a mut DynHttpService,
192 host: &'a [u8],
193 headers: &'a HttpHeaders,
194 connection: Pin<&'a mut AsyncIoStream>,
195 response: Pin<&'a mut ConnectResponse>,
196 settings: HttpConnectSettings<'a>,
197 ) -> Result<()>;
198 }
199 
200 impl Box<DynHttpService> {}
201}
202 
203assert_eq_size!(ffi::HttpConnectSettings, [u8; 16]);
204assert_eq_align!(ffi::HttpConnectSettings, u64);
205 
206pub type HeaderId = ffi::BuiltinIndicesEnum;
207pub type HeaderTable = ffi::HttpHeaderTable;
208pub type CustomHeader = ffi::HttpHeaderId;
209 
210// TODO(tewaro) soon: replace by enum HeaderId
211 
212/// Non-owning reference to a `kj::HttpHeaderId`.
213///
214/// `CustomHeader` is an opaque CXX type representing `HttpHeaderId` and can only be passed by
215/// reference across the FFI boundary. This wrapper makes the borrow lifetime explicit and provides
216/// a safe Rust handle.
217///
218/// `repr(transparent)` guarantees the same layout as `&ffi::HttpHeaderId` (i.e. a single
219/// pointer), which allows safe reinterpretation of `&[*const HttpHeaderId]` slices received
220/// from C++ into `&[CustomHeaderId]` via [`CustomHeaderId::from_ptr_slice`].
221#[derive(Clone, Copy)]
222#[repr(transparent)]
223pub struct CustomHeaderId<'a>(&'a CustomHeader);
224 
225impl<'a> From<&'a CustomHeader> for CustomHeaderId<'a> {
226 fn from(value: &'a CustomHeader) -> Self {
227 CustomHeaderId(value)
228 }
229}
230 
231impl<'a> CustomHeaderId<'a> {
232 /// Reinterpret a slice of `*const CustomHeader` pointers (as received from C++ via CXX) into
233 /// a slice of `HttpHeader`.
234 ///
235 /// This is the canonical way to receive a `kj::ArrayPtr<const kj::HttpHeaderId>` from C++:
236 /// the C++ side converts the array into a `rust::Slice<const HttpHeaderId* const>` and the
237 /// Rust side calls this function to get a safe `&[HttpHeaderIdRef]`.
238 ///
239 /// # Safety
240 /// `CustomHeaderId` is `#[repr(transparent)]` over `&CustomHeader`, which has
241 /// the same layout as `*const kj::HttpHeaderId`. The caller guarantees all pointers are valid.
242 pub unsafe fn from_ptr_slice(ptrs: &'a [*const CustomHeader]) -> &'a [Self] {
243 let ptr = std::ptr::from_ref::<[*const CustomHeader]>(ptrs) as *const [Self];
244 // SAFETY: CustomHeaderId is #[repr(transparent)] over &CustomHeader; caller guarantees all pointers are valid.
245 unsafe { &*ptr }
246 }
247}
248 
249/// Non-owning constant reference to `kj::HttpHeaders`
250#[derive(Clone, Copy)]
251pub struct HeadersRef<'a>(&'a ffi::HttpHeaders);
252 
253impl HeadersRef<'_> {
254 pub fn get(&self, id: HeaderId) -> Option<&[u8]> {
255 // SAFETY: self.0 is a valid HttpHeaders reference and id is a valid builtin header enum.
256 unsafe { ffi::get_header(self.0, id).into() }
257 }
258 
259 /// Look up a header by its `kj::HttpHeaderId`. This works for both builtin headers and custom
260 /// headers registered via `HttpHeaderTable::Builder::add()`.
261 pub fn get_by_id(&self, id: CustomHeaderId<'_>) -> Option<&[u8]> {
262 // SAFETY: self.0 is a valid HttpHeaders reference and id.0 is a valid HttpHeaderId.
263 unsafe { ffi::get_header_by_id(self.0, id.0).into() }
264 }
265 
266 #[must_use]
267 pub fn clone_shallow(&self) -> Headers<'_> {
268 Headers {
269 own: ffi::clone_shallow(self.0),
270 _marker: PhantomData,
271 }
272 }
273}
274 
275impl<'a> From<&'a ffi::HttpHeaders> for HeadersRef<'a> {
276 fn from(value: &'a ffi::HttpHeaders) -> Self {
277 HeadersRef(value)
278 }
279}
280 
281/// `HttpHeaders` that `kj::Own` the underlying C++ header object.
282///
283/// Notice, that despite the fact that headers are fully owned, because of a `shallowClone`
284/// method, data might not be owned: hence the lifetime parameter.
285pub struct Headers<'a> {
286 own: KjOwn<ffi::HttpHeaders>,
287 _marker: PhantomData<&'a ffi::HttpHeaders>,
288}
289 
290impl<'a> Headers<'a> {
291 #[must_use]
292 pub fn new(table: &'a HeaderTable) -> Self {
293 Self {
294 own: ffi::new_http_headers(table),
295 _marker: PhantomData,
296 }
297 }
298 
299 pub fn set(&mut self, id: HeaderId, value: &str) {
300 ffi::set_header(self.own.as_mut(), id, value);
301 }
302 
303 pub fn as_ref(&'a self) -> HeadersRef<'a> {
304 HeadersRef(self.own.as_ref())
305 }
306}
307 
308impl<'a, 'b> From<&'b Headers<'a>> for HeadersRef<'b> {
309 fn from(value: &'b Headers<'a>) -> Self {
310 value.as_ref()
311 }
312}
313 
314pub type Method = ffi::HttpMethod;
315pub type ConnectSettings<'a> = ffi::HttpConnectSettings<'a>;
316 
317/// Non-owning mutable reference to `kj::HttpService::Response`.
318pub struct ServiceResponse<'a>(Pin<&'a mut ffi::HttpServiceResponse>);
319 
320impl<'a> ServiceResponse<'a> {
321 /// Send response metadata and obtain the writable response body stream.
322 pub fn send<'h>(
323 self,
324 status_code: u32,
325 status_text: &str,
326 headers: impl Into<HeadersRef<'h>>,
327 expected_body_size: Option<u64>,
328 ) -> Result<crate::io::AsyncOutputStream<'a>> {
329 Ok(ffi::response_send(
330 self.0,
331 status_code,
332 status_text,
333 headers.into().0,
334 expected_body_size.into(),
335 )?
336 .into())
337 }
338 
339 pub(crate) fn into_ffi(self) -> Pin<&'a mut ffi::HttpServiceResponse> {
340 self.0
341 }
342}
343 
344impl<'a> From<Pin<&'a mut ffi::HttpServiceResponse>> for ServiceResponse<'a> {
345 fn from(value: Pin<&'a mut ffi::HttpServiceResponse>) -> Self {
346 Self(value)
347 }
348}
349 
350/// Non-owning mutable reference to `kj::HttpService::ConnectResponse`.
351pub struct ConnectResponse<'a>(Pin<&'a mut ffi::ConnectResponse>);
352 
353impl<'a> ConnectResponse<'a> {
354 /// Accept the CONNECT request without a response body.
355 pub fn accept<'h>(
356 self,
357 status_code: u32,
358 status_text: &str,
359 headers: impl Into<HeadersRef<'h>>,
360 ) -> Result<()> {
361 Ok(ffi::connect_response_accept(
362 self.0,
363 status_code,
364 status_text,
365 headers.into().0,
366 )?)
367 }
368 
369 /// Reject the CONNECT request and obtain the writable rejection body stream.
370 pub fn reject<'h>(
371 self,
372 status_code: u32,
373 status_text: &str,
374 headers: impl Into<HeadersRef<'h>>,
375 expected_body_size: Option<u64>,
376 ) -> Result<crate::io::AsyncOutputStream<'a>> {
377 Ok(ffi::connect_response_reject(
378 self.0,
379 status_code,
380 status_text,
381 headers.into().0,
382 expected_body_size.into(),
383 )?
384 .into())
385 }
386 
387 pub(crate) fn into_ffi(self) -> Pin<&'a mut ffi::ConnectResponse> {
388 self.0
389 }
390}
391 
392impl<'a> From<Pin<&'a mut ffi::ConnectResponse>> for ConnectResponse<'a> {
393 fn from(value: Pin<&'a mut ffi::ConnectResponse>) -> Self {
394 Self(value)
395 }
396}
397 
398#[async_trait::async_trait(?Send)]
399pub trait Service {
400 /// Make an HTTP request.
401 async fn request<'a>(
402 &'a mut self,
403 method: Method,
404 url: &'a [u8],
405 headers: HeadersRef<'a>,
406 request_body: Pin<&'a mut AsyncInputStream>,
407 response: ServiceResponse<'a>,
408 ) -> Result<()>;
409 
410 /// Make a CONNECT request
411 ///
412 /// WARNING: as c++ implementation does, this method has an outbound _immediate_
413 /// parameter inside `settings` (`tls_starter`).
414 ///
415 /// This method should be implement without using `async` method, since its body is called only
416 /// when it is first polled, but by manually creating a future using async block.
417 async fn connect<'a>(
418 &'a mut self,
419 host: &'a [u8],
420 headers: HeadersRef<'a>,
421 connection: Pin<&'a mut AsyncIoStream>,
422 response: ConnectResponse<'a>,
423 settings: ConnectSettings<'a>,
424 ) -> Result<()>;
425 
426 /// Convert `self` into the structure suitable for passing to C++ through FFI layer
427 /// To obtain `kj::HttpService` on C++ side finish wrapping with `fromRust()`
428 fn into_ffi(self) -> Box<DynHttpService>
429 where
430 Self: Sized + 'static,
431 {
432 Box::new(DynHttpService(Box::new(self)))
433 }
434}
435 
436pub struct CxxService<'a>(OwnOrMut<'a, ffi::HttpService>);
437 
438#[async_trait::async_trait(?Send)]
439impl Service for CxxService<'_> {
440 async fn request<'a>(
441 &'a mut self,
442 method: Method,
443 url: &'a [u8],
444 headers: HeadersRef<'a>,
445 request_body: Pin<&'a mut AsyncInputStream>,
446 response: ServiceResponse<'a>,
447 ) -> Result<()> {
448 let service = self.0.as_mut();
449 ffi::request(
450 service,
451 method,
452 url,
453 headers.0,
454 request_body,
455 response.into_ffi(),
456 )
457 .await?;
458 Ok(())
459 }
460 
461 fn connect<'a, 'b>(
462 &'a mut self,
463 host: &'a [u8],
464 headers: HeadersRef<'a>,
465 connection: Pin<&'a mut AsyncIoStream>,
466 response: ConnectResponse<'a>,
467 settings: ConnectSettings<'a>,
468 ) -> ::core::pin::Pin<Box<dyn ::core::future::Future<Output = Result<()>> + 'b>>
469 where
470 'a: 'b,
471 Self: 'b,
472 {
473 let service = self.0.as_mut();
474 Box::pin(
475 ffi::connect(
476 service,
477 host,
478 headers.0,
479 connection,
480 response.into_ffi(),
481 settings,
482 )
483 .map_err(Into::into),
484 )
485 }
486}
487 
488impl From<KjOwn<ffi::HttpService>> for CxxService<'_> {
489 fn from(value: KjOwn<ffi::HttpService>) -> Self {
490 CxxService(value.into())
491 }
492}
493 
494pub struct DynHttpService(Box<dyn Service>);
495 
496impl DynHttpService {
497 async fn request<'a>(
498 &'a mut self,
499 method: Method,
500 url: &'a [u8],
501 headers: &'a ffi::HttpHeaders,
502 request_body: Pin<&'a mut AsyncInputStream>,
503 response: Pin<&'a mut ffi::HttpServiceResponse>,
504 ) -> Result<()> {
505 let response = ServiceResponse::from(response);
506 self.0
507 .request(method, url, HeadersRef(headers), request_body, response)
508 .await?;
509 Ok(())
510 }
511 
512 fn connect<'a>(
513 &'a mut self,
514 host: &'a [u8],
515 headers: &'a ffi::HttpHeaders,
516 connection: Pin<&'a mut AsyncIoStream>,
517 response: Pin<&'a mut ffi::ConnectResponse>,
518 settings: ConnectSettings<'a>,
519 ) -> impl Future<Output = Result<()>> {
520 let headers = HeadersRef(headers);
521 let response = ConnectResponse::from(response);
522 self.0
523 .connect(host, headers, connection, response, settings)
524 }
525}