Skip to main content

rvlib/
httpserver.rs

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        // since the requests will only change the shown image, they will be handled sequentially
70        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                // we don't care about headers and bodies and strip everything after the first CRLF
76                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            // send recieved path
91            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}