use crate::parser::Http1Parser;
use crate::types::{Http1Config, Http1Error, HttpRequest};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum Http1ConnectionState {
Waiting,
ReadingRequest,
ReadingHeaders,
ReadingBody,
Processing,
SendingResponse,
Closing,
Closed,
}
#[derive(Debug)]
pub struct Http1Connection {
state: Http1ConnectionState,
parser: Http1Parser,
current_request: Option<HttpRequest>,
keep_alive: bool,
requests_handled: u64,
closed: bool,
last_activity_ms: u64,
time_initialized: bool,
}
impl Http1Connection {
#[inline]
pub fn new(config: Http1Config) -> Self {
Self {
keep_alive: config.keep_alive,
parser: Http1Parser::new(config),
state: Http1ConnectionState::Waiting,
current_request: None,
requests_handled: 0,
closed: false,
last_activity_ms: 0,
time_initialized: false,
}
}
#[inline]
pub fn state(&self) -> Http1ConnectionState {
self.state
}
#[inline]
pub fn is_closed(&self) -> bool {
self.closed
}
pub fn on_data(&mut self, data: &[u8]) -> Result<(Option<HttpRequest>, usize), Http1Error> {
if self.closed {
return Err(Http1Error::ConnectionClosed);
}
if !data.is_empty() && self.time_initialized {
self.parser.note_activity(self.last_activity_ms);
}
let (req, consumed) = match self.parser.feed(data) {
Ok(v) => v,
Err(e) => {
self.on_error(&e);
return Err(e);
}
};
if let Some(mut req) = req {
if !req.keep_alive {
self.keep_alive = false;
}
if req.chunked {
self.state = Http1ConnectionState::Processing;
} else if let Some(len) = req.content_length {
if len == 0 {
self.state = Http1ConnectionState::Processing;
} else if !req.body.is_empty() && req.body.len() as u64 == len {
self.state = Http1ConnectionState::Processing;
} else {
self.state = Http1ConnectionState::ReadingBody;
}
} else if req.line.method.as_ref() == "POST" || req.line.method.as_ref() == "PUT" {
req.content_length = Some(0);
self.state = Http1ConnectionState::Processing;
} else {
self.state = Http1ConnectionState::Processing;
}
self.current_request = Some(req);
self.requests_handled += 1;
Ok((self.current_request.clone(), consumed))
} else {
self.state = match self.parser.state() {
crate::parser::ParserState::WaitingRequest
| crate::parser::ParserState::ReadingRequest => Http1ConnectionState::ReadingRequest,
crate::parser::ParserState::ReadingHeaders => Http1ConnectionState::ReadingHeaders,
crate::parser::ParserState::ReadingBody
| crate::parser::ParserState::ReadingChunkSize
| crate::parser::ParserState::ReadingChunkData
| crate::parser::ParserState::ReadingChunkTrailer => Http1ConnectionState::ReadingBody,
crate::parser::ParserState::HeadersComplete => Http1ConnectionState::Processing,
crate::parser::ParserState::Error => Http1ConnectionState::Closed,
};
Ok((None, consumed))
}
}
pub fn body_complete(&mut self, body: Vec<u8>) {
if let Some(req) = self.current_request.as_mut() {
req.body = body;
self.state = Http1ConnectionState::Processing;
}
}
#[inline]
pub fn current_request(&self) -> Option<&HttpRequest> {
self.current_request.as_ref()
}
pub fn response_sent(&mut self) {
self.current_request = None;
if !self.keep_alive {
self.state = Http1ConnectionState::Closing;
self.closed = true;
} else {
self.parser.reset();
self.state = Http1ConnectionState::Waiting;
}
}
pub fn on_error(&mut self, _err: &Http1Error) {
self.state = Http1ConnectionState::Closing;
self.closed = true;
}
#[inline]
pub fn requests_handled(&self) -> u64 {
self.requests_handled
}
#[inline]
pub fn update_time(&mut self, now_ms: u64) {
self.last_activity_ms = now_ms;
self.time_initialized = true;
}
pub fn check_timeout(&mut self) -> Result<(), Http1Error> {
if self.closed {
return Ok(());
}
if let Err(e) = self.parser.check_idle_timeout(self.last_activity_ms) {
self.on_error(&e);
return Err(e);
}
Ok(())
}
pub fn close(&mut self) {
self.state = Http1ConnectionState::Closed;
self.closed = true;
}
}
impl Default for Http1Connection {
fn default() -> Self {
Self::new(Http1Config::new())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_connection_basic() {
let mut conn = Http1Connection::new(Http1Config::new());
let input = b"GET / HTTP/1.1\r\nHost: example.com\r\n\r\n";
let (req, _) = conn.on_data(input).unwrap();
assert!(req.is_some());
assert_eq!(conn.state(), Http1ConnectionState::Processing);
}
#[test]
fn test_connection_lifecycle() {
let mut conn = Http1Connection::new(Http1Config::new());
let input = b"GET / HTTP/1.1\r\nHost: example.com\r\n\r\n";
conn.on_data(input).unwrap();
assert_eq!(conn.state(), Http1ConnectionState::Processing);
conn.response_sent();
assert_eq!(conn.state(), Http1ConnectionState::Waiting);
assert!(!conn.is_closed());
}
#[test]
fn test_connection_close() {
let mut conn = Http1Connection::new(Http1Config::new());
let input =
b"GET / HTTP/1.1\r\nHost: example.com\r\nConnection: close\r\n\r\n";
conn.on_data(input).unwrap();
conn.response_sent();
assert!(conn.is_closed());
}
#[test]
fn test_connection_error_cases() {
let mut conn = Http1Connection::new(Http1Config::new());
let r = conn.on_data(b"INV HTTP/1.1\r\nHost: x\r\n\r\n");
assert!(r.is_err());
assert!(conn.is_closed());
}
#[test]
fn test_connection_body_post() {
let mut conn = Http1Connection::new(Http1Config::new());
let input =
b"POST /api HTTP/1.1\r\nHost: example.com\r\nContent-Length: 5\r\n\r\nhello";
let (req, _) = conn.on_data(input).unwrap();
let req = req.unwrap();
assert_eq!(req.content_length, Some(5));
assert_eq!(req.body, b"hello");
assert_eq!(conn.state(), Http1ConnectionState::Processing);
}
#[test]
fn test_connection_state_variants() {
let states = [
Http1ConnectionState::Waiting,
Http1ConnectionState::ReadingRequest,
Http1ConnectionState::ReadingHeaders,
Http1ConnectionState::ReadingBody,
Http1ConnectionState::Processing,
Http1ConnectionState::SendingResponse,
Http1ConnectionState::Closing,
Http1ConnectionState::Closed,
];
for (i, s) in states.iter().enumerate() {
assert_eq!(*s, states[i]);
}
assert_ne!(Http1ConnectionState::Waiting, Http1ConnectionState::Closed);
}
#[test]
fn test_connection_default() {
let conn = Http1Connection::default();
assert_eq!(conn.state(), Http1ConnectionState::Waiting);
assert!(!conn.is_closed());
assert_eq!(conn.requests_handled(), 0);
}
#[test]
fn test_connection_keep_alive_multiple_requests() {
let mut conn = Http1Connection::new(Http1Config::new());
let input1 = b"GET /1 HTTP/1.1\r\nHost: example.com\r\n\r\n";
let (req1, _) = conn.on_data(input1).unwrap();
assert!(req1.is_some());
assert_eq!(conn.requests_handled(), 1);
conn.response_sent();
assert_eq!(conn.state(), Http1ConnectionState::Waiting);
let input2 = b"GET /2 HTTP/1.1\r\nHost: example.com\r\n\r\n";
let (req2, _) = conn.on_data(input2).unwrap();
assert!(req2.is_some());
assert_eq!(conn.requests_handled(), 2);
}
#[test]
fn test_connection_closed_rejects_data() {
let mut conn = Http1Connection::new(Http1Config::new());
conn.close();
assert!(conn.is_closed());
assert_eq!(conn.state(), Http1ConnectionState::Closed);
let input = b"GET / HTTP/1.1\r\nHost: example.com\r\n\r\n";
let r = conn.on_data(input);
assert!(r.is_err());
assert!(matches!(r.unwrap_err(), Http1Error::ConnectionClosed));
}
#[test]
fn test_connection_current_request() {
let mut conn = Http1Connection::new(Http1Config::new());
assert!(conn.current_request().is_none());
let input = b"GET / HTTP/1.1\r\nHost: example.com\r\n\r\n";
conn.on_data(input).unwrap();
assert!(conn.current_request().is_some());
assert_eq!(conn.current_request().unwrap().line.method.as_ref(), "GET");
}
#[test]
fn test_connection_chunked_body_state() {
let mut conn = Http1Connection::new(Http1Config::new());
let input = b"POST /upload HTTP/1.1\r\nHost: example.com\r\nTransfer-Encoding: chunked\r\n\r\n5\r\nHello\r\n0\r\n\r\n";
let (req, _) = conn.on_data(input).unwrap();
assert!(req.is_some());
let req = req.unwrap();
assert_eq!(conn.state(), Http1ConnectionState::Processing);
assert_eq!(req.body, b"Hello");
}
#[test]
fn test_connection_chunked_body_partial() {
let mut conn = Http1Connection::new(Http1Config::new());
let part1 = b"POST /upload HTTP/1.1\r\nHost: example.com\r\nTransfer-Encoding: chunked\r\n\r\n";
let (req, _) = conn.on_data(part1).unwrap();
assert!(req.is_none());
assert_eq!(conn.state(), Http1ConnectionState::ReadingBody);
let part2 = b"5\r\nHello\r\n0\r\n\r\n";
let (req, _) = conn.on_data(part2).unwrap();
assert!(req.is_some());
assert_eq!(conn.state(), Http1ConnectionState::Processing);
assert_eq!(req.unwrap().body, b"Hello");
}
#[test]
fn test_connection_zero_content_length() {
let mut conn = Http1Connection::new(Http1Config::new());
let input = b"POST /api HTTP/1.1\r\nHost: example.com\r\nContent-Length: 0\r\n\r\n";
let (req, _) = conn.on_data(input).unwrap();
assert!(req.is_some());
assert_eq!(conn.state(), Http1ConnectionState::Processing);
}
#[test]
fn test_connection_put_method() {
let mut conn = Http1Connection::new(Http1Config::new());
let input =
b"PUT /resource HTTP/1.1\r\nHost: example.com\r\nContent-Length: 5\r\n\r\nworld";
let (req, _) = conn.on_data(input).unwrap();
let req = req.unwrap();
assert_eq!(req.line.method.as_ref(), "PUT");
assert_eq!(req.body, b"world");
assert_eq!(conn.state(), Http1ConnectionState::Processing);
}
#[test]
fn test_connection_head_method() {
let mut conn = Http1Connection::new(Http1Config::new());
let input = b"HEAD / HTTP/1.1\r\nHost: example.com\r\n\r\n";
let (req, _) = conn.on_data(input).unwrap();
assert!(req.is_some());
assert_eq!(req.unwrap().line.method.as_ref(), "HEAD");
assert_eq!(conn.state(), Http1ConnectionState::Processing);
}
#[test]
fn test_connection_on_error_closes_connection() {
let mut conn = Http1Connection::new(Http1Config::new());
conn.on_error(&Http1Error::SyntaxError("test".into()));
assert!(conn.is_closed());
assert_eq!(conn.state(), Http1ConnectionState::Closing);
}
#[test]
fn test_connection_body_complete_no_current_request() {
let mut conn = Http1Connection::new(Http1Config::new());
conn.body_complete(b"test".to_vec());
assert_eq!(conn.state(), Http1ConnectionState::Waiting);
}
#[test]
fn test_connection_response_sent_clears_request() {
let mut conn = Http1Connection::new(Http1Config::new());
let input = b"GET / HTTP/1.1\r\nHost: example.com\r\n\r\n";
conn.on_data(input).unwrap();
assert!(conn.current_request().is_some());
conn.response_sent();
assert!(conn.current_request().is_none());
}
#[test]
fn test_connection_closing_state() {
let mut conn = Http1Connection::new(Http1Config::new());
let input = b"GET / HTTP/1.1\r\nHost: example.com\r\nConnection: close\r\n\r\n";
conn.on_data(input).unwrap();
conn.response_sent();
assert_eq!(conn.state(), Http1ConnectionState::Closing);
assert!(conn.is_closed());
}
#[test]
fn test_connection_debug_format() {
let conn = Http1Connection::new(Http1Config::new());
let s = format!("{:?}", conn);
assert!(!s.is_empty());
}
#[test]
fn test_connection_idle_timeout_during_body_read() {
let config = Http1Config::new().with_idle_timeout_ms(30_000);
let mut conn = Http1Connection::new(config);
conn.update_time(0);
let (req, _) = conn
.on_data(b"POST / HTTP/1.1\r\nHost: example.com\r\nContent-Length: 10\r\n\r\n")
.unwrap();
assert!(req.is_none());
assert_eq!(conn.state(), Http1ConnectionState::ReadingBody);
conn.update_time(30_001);
let r = conn.check_timeout();
assert!(matches!(r, Err(Http1Error::IdleTimeout)), "实际 {r:?}");
assert!(conn.is_closed());
}
#[test]
fn test_connection_body_activity_prevents_timeout() {
let config = Http1Config::new().with_idle_timeout_ms(30_000);
let mut conn = Http1Connection::new(config);
conn.update_time(0);
let (req, _) = conn
.on_data(b"POST / HTTP/1.1\r\nHost: example.com\r\nContent-Length: 6\r\n\r\nhe")
.unwrap();
assert!(req.is_none());
conn.update_time(20_000);
let (req, _) = conn.on_data(b"ll").unwrap();
assert!(req.is_none());
conn.update_time(45_000);
assert!(conn.check_timeout().is_ok(), "活跃刷新后不得超时");
let (req, _) = conn.on_data(b"lo").unwrap();
assert!(req.is_some());
}
}