1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
use crate::blocking::{runtime, Body, BodyShim, BodyStreamer, Client, Response};
use crate::Request;
use conjure_error::Error;
use conjure_object::BearerToken;
use futures::channel::oneshot;
use futures::executor;
use hyper::{HeaderMap, Method};
pub struct RequestBuilder<'a> {
client: &'a Client,
request: Request<'static>,
streamer: Option<BodyStreamer<Box<dyn Body + 'a>>>,
}
impl<'a> RequestBuilder<'a> {
pub(crate) fn new(
client: &'a Client,
method: Method,
pattern: &'static str,
) -> RequestBuilder<'a> {
RequestBuilder {
client,
request: Request::new(&client.0, method, pattern),
streamer: None,
}
}
pub fn headers_mut(&mut self) -> &mut HeaderMap {
&mut self.request.headers
}
pub fn bearer_token(mut self, token: &BearerToken) -> RequestBuilder<'a> {
self.request.bearer_token(token);
self
}
#[allow(clippy::needless_pass_by_value)]
pub fn param<T>(mut self, name: &str, value: T) -> RequestBuilder<'a>
where
T: ToString,
{
self.request.param(name, value);
self
}
pub fn idempotent(mut self, idempotent: bool) -> RequestBuilder<'a> {
self.request.idempotent = idempotent;
self
}
pub fn body<T>(mut self, body: T) -> RequestBuilder<'a>
where
T: Body + 'a,
{
let (body, streamer) = BodyShim::new(Box::new(body) as _);
self.request.body(body);
self.streamer = Some(streamer);
self
}
pub fn send(self) -> Result<Response, Error> {
let (sender, receiver) = oneshot::channel();
let client = self.client.0.clone();
let request = self.request;
runtime().map_err(Error::internal_safe)?.spawn(async move {
let r = client.send(request).await;
let _ = sender.send(r);
});
if let Some(streamer) = self.streamer {
streamer.stream();
}
match executor::block_on(receiver) {
Ok(Ok(r)) => Ok(Response::new(r)),
Ok(Err(e)) => Err(e.with_backtrace()),
Err(e) => Err(Error::internal_safe(e)),
}
}
}