1use httparse::{EMPTY_HEADER, Request};
2use rvimage_domain::{RvResult, rverr, to_rv};
3use std::{
4 fmt::Debug,
5 io::{Read, prelude::*},
6 net::TcpListener,
7 str,
8 sync::mpsc::{self, Receiver, Sender},
9 thread::{self, JoinHandle},
10};
11use tracing::info;
12
13#[derive(Debug, PartialEq)]
14enum HandleResult {
15 Path(String),
16 Terminate,
17}
18
19fn handle_connection(buffer: &[u8]) -> RvResult<HandleResult> {
20 let mut headers = [EMPTY_HEADER; 64];
21
22 let mut req = Request::new(&mut headers);
23 let res = req.parse(buffer);
24 res.map_err(to_rv)?;
25 match req.path {
26 Some(p) => Ok(if p == "/TERMINATE" {
27 HandleResult::Terminate
28 } else {
29 HandleResult::Path(
30 percent_encoding::percent_decode_str(&p[1..])
31 .decode_utf8()
32 .map_err(to_rv)?
33 .to_string(),
34 )
35 }),
36 None => {
37 let msg = "could not find path in stream";
38 match str::from_utf8(buffer) {
39 Ok(b) => Err(rverr!("{} '{}'", msg, b)),
40 Err(e) => Err(rverr!("{}, {:?}", msg, e)),
41 }
42 }
43 }
44}
45
46fn to_rv_or_send<T, E>(tx: &Sender<RvResult<String>>, x: Result<T, E>) -> RvResult<T>
47where
48 E: Debug,
49{
50 match x {
51 Ok(r) => Ok(r),
52 Err(e) => {
53 let error_str = format!("{e:?}");
54 match tx.send(Err(to_rv(e))) {
55 Ok(()) => Err(rverr!("error in http server, {}", error_str)),
56 Err(e) => Err(rverr!("error {}, send error {:?}", error_str, e)),
57 }
58 }
59 }
60}
61pub type LaunchResultType = RvResult<(JoinHandle<RvResult<()>>, Receiver<RvResult<String>>)>;
62pub fn launch(address: String) -> LaunchResultType {
63 info!("spawning httpserver at {address}");
64 let (tx_from_server, rx_from_server) = mpsc::channel();
65 let handle = thread::spawn(move || -> RvResult<()> {
66 let bind_result = TcpListener::bind(address);
67 let listener = to_rv_or_send(&tx_from_server, bind_result)?;
68 let mut buffer = vec![0; 4096];
69 for stream in listener.incoming() {
71 #[allow(clippy::indexing_slicing)]
72 let stream_processing_result: RvResult<HandleResult> = {
73 let mut stream = stream.map_err(to_rv)?;
74 stream.read(&mut buffer).map_err(to_rv)?;
75 let buffer_slice = if let Some(idx) =
77 (0..(buffer.len() - 1)).find(|idx| &buffer[*idx..*idx + 2] == b"\r\n")
78 {
79 &buffer[..idx + 2]
80 } else {
81 &buffer
82 };
83 let response = "HTTP/1.1 200 OK\r\n\r\n";
84 stream.write(response.as_bytes()).map_err(to_rv)?;
85 stream.flush().map_err(to_rv)?;
86 let res = handle_connection(buffer_slice);
87 let path = res?;
88 Ok(path)
89 };
90 if let Ok(p) = to_rv_or_send(&tx_from_server, stream_processing_result) {
92 match p {
93 HandleResult::Terminate => {
94 info!("terminating httpserver");
95 return Ok(());
96 }
97 HandleResult::Path(p_) => {
98 info!("tcp listener sending result...");
99 let send_result = tx_from_server.send(Ok(p_));
100 info!("done. {send_result:?}");
101 to_rv_or_send(&tx_from_server, send_result)?;
102 }
103 }
104 }
105 info!("tcp listener waiting for new input");
106 }
107 Ok(())
108 });
109 info!("...done");
110 Ok((handle, rx_from_server))
111}
112
113fn increase_port(address: &str) -> RvResult<String> {
114 let address_wo_port = address.split(':').next();
115 let port = address.split(':').next_back();
116 if let Some(port) = port {
117 if let Some(address_wo_port) = address_wo_port {
118 Ok(format!(
119 "{}:{}",
120 address_wo_port,
121 (port.parse::<usize>().map_err(to_rv)? + 1)
122 ))
123 } else {
124 Err(rverr!("is address of {} missing?", address))
125 }
126 } else {
127 Err(rverr!("is port of address {} missing?", address))
128 }
129}
130
131pub fn restart_with_increased_port(
132 http_addr: &str,
133) -> RvResult<(String, Option<Receiver<RvResult<String>>>)> {
134 let http_addr = increase_port(http_addr)?;
135
136 info!("restarting http server with increased port");
137 Ok(if let Ok((_, rx)) = launch(http_addr.clone()) {
138 (http_addr, Some(rx))
139 } else {
140 (http_addr, None)
141 })
142}
143#[test]
144fn test_handler() -> RvResult<()> {
145 let buffer = b"garbage";
146 assert!(handle_connection(buffer.as_slice()).is_err());
147
148 let buffer = b"GET /index.html HTTP/1.1\r\nHost:";
149 assert_eq!(
150 handle_connection(buffer.as_slice()),
151 Ok(HandleResult::Path("index.html".to_string()))
152 );
153
154 let buffer = b"GET /folder%20name/file%20name.png HTTP/1.1\r\nHost:";
155 assert_eq!(
156 handle_connection(buffer.as_slice()),
157 Ok(HandleResult::Path("folder name/file name.png".to_string()))
158 );
159 let buffer = b"GET /TERMINATE HTTP/1.1\r\nHost:";
160 assert_eq!(
161 handle_connection(buffer.as_slice()),
162 Ok(HandleResult::Terminate)
163 );
164
165 Ok(())
166}
167#[cfg(test)]
168use std::{net::TcpStream, time::Duration};
169#[test]
170fn test_launch() -> RvResult<()> {
171 let address = "127.0.0.1:7942";
172 println!("launching server...");
173 let (handle, rx) = launch(address.to_string())?;
174 thread::sleep(Duration::from_millis(10));
175 assert!(!handle.is_finished());
176 println!("...done");
177
178 let send_request = |req| -> RvResult<()> {
179 let mut stream = TcpStream::connect(address).map_err(to_rv)?;
180 stream.write(req).map_err(to_rv)?;
181 stream.flush().map_err(to_rv)?;
182 Ok(())
183 };
184 println!("writing to stream...");
185 let input_stream = b"GET /some_path.png HTTP/1.1\r\nHost:";
186 send_request(input_stream.as_slice())?;
187 println!("...done");
188 println!("writing to stream...");
189 let input_stream = b"GET /some_other_path.png HTTP/1.1\r\nHost:";
190 send_request(input_stream.as_slice())?;
191 println!("...done");
192 thread::sleep(Duration::from_millis(1500));
193 println!("checking results...");
194 let result1 = rx.recv().map_err(to_rv)?;
195 let result2 = rx.recv().map_err(to_rv)?;
196 assert_eq!(result1, Ok("some_path.png".to_string()));
197 assert_eq!(result2, Ok("some_other_path.png".to_string()));
198 println!("...done");
199 println!("terminate...");
200 let terminate_stream = b"GET /TERMINATE HTTP/1.1\r\n";
201 send_request(terminate_stream.as_slice())?;
202 thread::sleep(Duration::from_millis(500));
203 assert!(handle.is_finished());
204 println!("...done");
205 Ok(())
206}
207#[test]
208fn test_increase_port() -> RvResult<()> {
209 assert_eq!(increase_port("address:1234")?, "address:1235");
210 Ok(())
211}