Skip to main content

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}