iceberg_catalog_rest/request.rs
1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements. See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership. The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License. You may obtain a copy of the License at
8//
9// http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied. See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! The request type REST catalog authentication works with.
19
20use http::{HeaderMap, Method};
21use iceberg::Result;
22#[cfg(feature = "sigv4")]
23use reqwest::Url;
24use reqwest::{Request, RequestBuilder};
25
26/// An outgoing REST request being authenticated by an
27/// [`AuthSession`](crate::AuthSession).
28///
29/// Wraps the request so a session mutates it through the stable
30/// `http` crate types rather than the concrete request type the REST catalog
31/// uses internally.
32pub struct HttpRequest {
33 inner: Request,
34}
35
36impl HttpRequest {
37 /// Wraps a request, e.g. to unit-test a custom
38 /// [`AuthSession`](crate::AuthSession).
39 pub fn new(inner: Request) -> Self {
40 Self { inner }
41 }
42
43 /// Builds the request `builder` describes.
44 pub(crate) fn build(builder: RequestBuilder) -> Result<Self> {
45 Ok(Self::new(builder.build()?))
46 }
47
48 /// The wrapped request, for the client that sends it.
49 pub(crate) fn into_inner(self) -> Request {
50 self.inner
51 }
52
53 /// The request URL.
54 #[cfg(feature = "sigv4")]
55 pub(crate) fn url(&self) -> &Url {
56 self.inner.url()
57 }
58
59 /// The mutable request URL, for signers that rewrite it.
60 #[cfg(feature = "sigv4")]
61 pub(crate) fn url_mut(&mut self) -> &mut Url {
62 self.inner.url_mut()
63 }
64
65 /// The request method.
66 pub fn method(&self) -> &Method {
67 self.inner.method()
68 }
69
70 /// The request URL, as a string (scheme, host, path and query).
71 pub fn url_str(&self) -> &str {
72 self.inner.url().as_str()
73 }
74
75 /// The request headers.
76 pub fn headers(&self) -> &HeaderMap {
77 self.inner.headers()
78 }
79
80 /// The mutable request headers, e.g. to add an `Authorization` header.
81 pub fn headers_mut(&mut self) -> &mut HeaderMap {
82 self.inner.headers_mut()
83 }
84
85 /// The request body, distinguishing an absent body from a streaming one:
86 /// signers can sign [`HttpRequestBody::Empty`] (empty-payload hash) and
87 /// [`HttpRequestBody::Buffered`], but not [`HttpRequestBody::Streaming`].
88 pub fn body(&self) -> HttpRequestBody<'_> {
89 match self.inner.body() {
90 None => HttpRequestBody::Empty,
91 Some(body) => match body.as_bytes() {
92 Some(bytes) => HttpRequestBody::Buffered(bytes),
93 None => HttpRequestBody::Streaming,
94 },
95 }
96 }
97}
98
99/// The body of an [`HttpRequest`], as seen by authentication.
100#[derive(Debug, Clone, Copy, PartialEq, Eq)]
101pub enum HttpRequestBody<'a> {
102 /// No body is set.
103 Empty,
104 /// An in-memory body.
105 Buffered(&'a [u8]),
106 /// A streaming body, whose bytes are not available for e.g. signing.
107 Streaming,
108}
109
110impl<'a> HttpRequestBody<'a> {
111 /// The signable bytes: empty for [`Self::Empty`], the buffer for
112 /// [`Self::Buffered`], and `None` for [`Self::Streaming`].
113 pub fn as_bytes(&self) -> Option<&'a [u8]> {
114 match self {
115 HttpRequestBody::Empty => Some(&[]),
116 HttpRequestBody::Buffered(bytes) => Some(bytes),
117 HttpRequestBody::Streaming => None,
118 }
119 }
120}
121
122#[cfg(test)]
123mod tests {
124 use reqwest::Client;
125
126 use super::*;
127
128 #[test]
129 fn test_http_request_body_states() {
130 let client = Client::new();
131
132 // No body at all.
133 let req = client
134 .get("https://rest.example.com/v1/config")
135 .build()
136 .unwrap();
137 let http_req = HttpRequest::new(req);
138 let body = http_req.body();
139 assert_eq!(body, HttpRequestBody::Empty);
140 assert_eq!(body.as_bytes(), Some(&[] as &[u8]));
141
142 // An in-memory body.
143 let req = client
144 .post("https://rest.example.com/v1/namespaces")
145 .body("{}")
146 .build()
147 .unwrap();
148 let http_req = HttpRequest::new(req);
149 let body = http_req.body();
150 assert_eq!(body, HttpRequestBody::Buffered(b"{}"));
151 assert_eq!(body.as_bytes(), Some(b"{}" as &[u8]));
152
153 // A streaming body: bytes are unavailable, so it must not sign as empty.
154 let req = client
155 .post("https://rest.example.com/v1/namespaces")
156 .body(reqwest::Body::wrap_stream(futures::stream::once(async {
157 Ok::<_, std::io::Error>(bytes::Bytes::from_static(b"chunk"))
158 })))
159 .build()
160 .unwrap();
161 let http_req = HttpRequest::new(req);
162 let body = http_req.body();
163 assert_eq!(body, HttpRequestBody::Streaming);
164 assert_eq!(body.as_bytes(), None);
165 }
166}