diff --git a/crates/solstone-linux/src/private_file.rs b/crates/solstone-linux/src/private_file.rs index 155b3c7..dcaf385 100644 --- a/crates/solstone-linux/src/private_file.rs +++ b/crates/solstone-linux/src/private_file.rs @@ -3,10 +3,10 @@ use std::{ fmt, - fs::{self, File}, + fs::File, io::{self, Write}, os::fd::AsFd, - path::Path, + path::{Component, Path}, sync::atomic::{AtomicU64, Ordering}, }; @@ -75,60 +75,78 @@ pub(crate) fn ensure_private_directory(path: &Path) -> Result<(), PrivateFileErr if path.as_os_str().is_empty() || path.parent().is_none() { return Err(PrivateFileError::InvalidTarget("directory")); } - let mut missing = Vec::new(); - let mut current = path; - loop { - match fs::symlink_metadata(current) { - Ok(metadata) if metadata.file_type().is_symlink() || !metadata.is_dir() => { - return Err(PrivateFileError::InvalidTarget("directory")); - } - Ok(_) => break, - Err(error) if error.kind() == io::ErrorKind::NotFound => { - missing.push(current.to_path_buf()); - current = current - .parent() - .filter(|parent| !parent.as_os_str().is_empty()) - .ok_or(PrivateFileError::InvalidTarget("directory"))?; - } - Err(error) => return Err(PrivateFileError::io("directory", "inspect", error)), + let mut current = open_walk_start(path.is_absolute())?; + let mut saw_component = false; + for component in path.components() { + let name = match component { + Component::RootDir | Component::CurDir => continue, + Component::ParentDir => Path::new("..").as_os_str(), + Component::Normal(name) => name, + Component::Prefix(_) => return Err(PrivateFileError::InvalidTarget("directory")), + }; + saw_component = true; + let (next, created) = open_or_create_directory(¤t, name)?; + if created { + set_and_verify_directory(&next)?; } + current = next; } - for directory in missing.iter().rev() { - fs::create_dir(directory) - .map_err(|error| PrivateFileError::io("directory", "create", error))?; - set_and_verify_mode(directory, 0o700, true)?; + if !saw_component { + return Err(PrivateFileError::InvalidTarget("directory")); } - set_and_verify_mode(path, 0o700, true) + set_and_verify_directory(¤t) } -fn set_and_verify_mode(path: &Path, mode: u32, directory: bool) -> Result<(), PrivateFileError> { - let flags = rustix::fs::OFlags::RDONLY - | rustix::fs::OFlags::CLOEXEC - | rustix::fs::OFlags::NOFOLLOW - | if directory { - rustix::fs::OFlags::DIRECTORY +fn open_walk_start(absolute: bool) -> Result { + open_directory_at( + rustix::fs::CWD, + if absolute { + Path::new("/") } else { - rustix::fs::OFlags::empty() - }; - let descriptor = rustix::fs::openat(rustix::fs::CWD, path, flags, rustix::fs::Mode::empty()) - .map_err(|error| { - if error == rustix::io::Errno::LOOP || error == rustix::io::Errno::NOTDIR { - PrivateFileError::InvalidTarget("target") - } else { - PrivateFileError::io("target", "open", error.into()) + Path::new(".") + }, + ) +} + +fn open_or_create_directory( + parent: &File, + name: &std::ffi::OsStr, +) -> Result<(File, bool), PrivateFileError> { + match open_directory_at(parent, Path::new(name)) { + Ok(directory) => Ok((directory, false)), + Err(PrivateFileError::Io { + kind: io::ErrorKind::NotFound, + .. + }) => { + match rustix::fs::mkdirat( + parent, + name, + rustix::fs::Mode::RUSR | rustix::fs::Mode::WUSR | rustix::fs::Mode::XUSR, + ) { + Ok(()) => {} + Err(rustix::io::Errno::EXIST) => { + return open_directory_at(parent, Path::new(name)) + .map(|directory| (directory, false)); + } + Err(error) => { + return Err(PrivateFileError::io("directory", "create", error.into())); + } } - })?; - let expected = rustix::fs::Mode::from_raw_mode(mode); - rustix::fs::fchmod(&descriptor, expected) + open_directory_at(parent, Path::new(name)).map(|directory| (directory, true)) + } + Err(error) => Err(error), + } +} + +fn set_and_verify_directory(descriptor: &File) -> Result<(), PrivateFileError> { + let expected = rustix::fs::Mode::RUSR | rustix::fs::Mode::WUSR | rustix::fs::Mode::XUSR; + rustix::fs::fchmod(descriptor, expected) .map_err(|error| PrivateFileError::io("target", "chmod", error.into()))?; - let stat = rustix::fs::fstat(&descriptor) + let stat = rustix::fs::fstat(descriptor) .map_err(|error| PrivateFileError::io("target", "inspect", error.into()))?; - let valid_kind = if directory { - rustix::fs::FileType::from_raw_mode(stat.st_mode) == rustix::fs::FileType::Directory - } else { - rustix::fs::FileType::from_raw_mode(stat.st_mode) == rustix::fs::FileType::RegularFile - }; - if !valid_kind || rustix::fs::Mode::from_raw_mode(stat.st_mode) != expected { + if rustix::fs::FileType::from_raw_mode(stat.st_mode) != rustix::fs::FileType::Directory + || rustix::fs::Mode::from_raw_mode(stat.st_mode) != expected + { return Err(PrivateFileError::InvalidTarget("target")); } Ok(()) @@ -198,22 +216,16 @@ pub(crate) fn atomic_write_bytes_with_fault( TEMP_COUNTER.fetch_add(1, Ordering::Relaxed), name = name.to_string_lossy() ); - let result = write_temporary(&parent_descriptor, name, &temporary, bytes, fault); - if result.is_err() { - let _ = rustix::fs::unlinkat( - &parent_descriptor, - temporary.as_str(), - rustix::fs::AtFlags::empty(), - ); - // Before rename the target is untouched; after rename it contains the new - // complete value. Rewriting either value after an error cannot be made safe. - } - result + write_temporary(&parent_descriptor, name, &temporary, bytes, fault) } fn open_directory(path: &Path) -> Result { + open_directory_at(rustix::fs::CWD, path) +} + +fn open_directory_at(parent: Fd, path: &Path) -> Result { rustix::fs::openat( - rustix::fs::CWD, + parent, path, rustix::fs::OFlags::RDONLY | rustix::fs::OFlags::CLOEXEC @@ -253,36 +265,49 @@ fn write_temporary( ) .map_err(|error| PrivateFileError::io("file", "create", error.into()))?; let mut file = File::from(descriptor); - rustix::fs::fchmod(&file, rustix::fs::Mode::RUSR | rustix::fs::Mode::WUSR) - .map_err(|error| PrivateFileError::io("file", "chmod", error.into()))?; - fault - .before(DurableWriteStage::Write) - .map_err(|error| PrivateFileError::io("file", "write", error))?; - file.write_all(bytes) - .and_then(|()| file.flush()) - .map_err(|error| PrivateFileError::io("file", "write", error))?; - fault - .before(DurableWriteStage::Fsync) - .map_err(|error| PrivateFileError::io("file", "fsync", error))?; - file.sync_all() - .map_err(|error| PrivateFileError::io("file", "fsync", error))?; - fault - .before(DurableWriteStage::Rename) - .map_err(|error| PrivateFileError::io("file", "rename", error))?; - rustix::fs::renameat(parent, temporary, parent, name) - .map_err(|error| PrivateFileError::io("file", "rename", error.into()))?; - fault - .before(DurableWriteStage::DirSync) - .map_err(|error| PrivateFileError::io("directory", "fsync", error))?; - parent - .sync_all() - .map_err(|error| PrivateFileError::io("directory", "fsync", error)) + let mut temporary_exists = true; + let result = (|| { + rustix::fs::fchmod(&file, rustix::fs::Mode::RUSR | rustix::fs::Mode::WUSR) + .map_err(|error| PrivateFileError::io("file", "chmod", error.into()))?; + fault + .before(DurableWriteStage::Write) + .map_err(|error| PrivateFileError::io("file", "write", error))?; + file.write_all(bytes) + .and_then(|()| file.flush()) + .map_err(|error| PrivateFileError::io("file", "write", error))?; + fault + .before(DurableWriteStage::Fsync) + .map_err(|error| PrivateFileError::io("file", "fsync", error))?; + file.sync_all() + .map_err(|error| PrivateFileError::io("file", "fsync", error))?; + fault + .before(DurableWriteStage::Rename) + .map_err(|error| PrivateFileError::io("file", "rename", error))?; + rustix::fs::renameat(parent, temporary, parent, name) + .map_err(|error| PrivateFileError::io("file", "rename", error.into()))?; + temporary_exists = false; + fault + .before(DurableWriteStage::DirSync) + .map_err(|error| PrivateFileError::io("directory", "fsync", error))?; + parent + .sync_all() + .map_err(|error| PrivateFileError::io("directory", "fsync", error)) + })(); + if result.is_err() && temporary_exists { + let _ = rustix::fs::unlinkat(parent, temporary, rustix::fs::AtFlags::empty()); + // Before rename the target is untouched; after rename it contains the new + // complete value. Rewriting either value after an error cannot be made safe. + } + result } #[cfg(test)] mod tests { use super::*; - use std::os::unix::fs::{MetadataExt, PermissionsExt, symlink}; + use std::{ + fs, + os::unix::fs::{MetadataExt, PermissionsExt, symlink}, + }; struct FailAt(DurableWriteStage); impl DurableWriteFault for FailAt { @@ -417,6 +442,34 @@ mod tests { } } + #[test] + fn create_collision_does_not_remove_an_unowned_temporary_file() { + let temp = tempfile::tempdir().unwrap(); + let parent = open_directory(temp.path()).unwrap(); + let temporary = ".state.collision.tmp"; + fs::write(temp.path().join(temporary), b"owned elsewhere").unwrap(); + let error = write_temporary( + &parent, + std::ffi::OsStr::new("state"), + temporary, + b"replacement", + &NoWriteFault, + ) + .unwrap_err(); + assert!(matches!( + error, + PrivateFileError::Io { + kind: io::ErrorKind::AlreadyExists, + .. + } + )); + assert_eq!( + fs::read(temp.path().join(temporary)).unwrap(), + b"owned elsewhere" + ); + assert!(!temp.path().join("state").exists()); + } + #[test] fn failed_initial_write_never_leaves_partial_target_or_temporary() { let temp = tempfile::tempdir().unwrap(); diff --git a/crates/solstone-linux/src/private_link.rs b/crates/solstone-linux/src/private_link.rs index bdc3c39..d119fcf 100644 --- a/crates/solstone-linux/src/private_link.rs +++ b/crates/solstone-linux/src/private_link.rs @@ -461,34 +461,52 @@ pub(crate) fn load_observer( }; let observer = serde_json::from_slice::(&bytes) .map_err(|_| PrivateStateError::MalformedObserver)?; - if observer.credential_instance_id != credential_instance_id - || observer.name != expected_name - || observer.protocol_version != 2 - || observer.key.is_empty() - || observer.prefix.is_empty() - || observer.name.is_empty() - || observer.ingest_url.is_empty() - || contains_invalid_header_value(&observer.key) - || confine_path(origin, &observer.ingest_url).is_err() - { + if !observer_is_valid(&observer, credential_instance_id, expected_name, origin) { return Ok(None); } Ok(Some(observer)) } +fn observer_is_valid( + observer: &ObserverState, + credential_instance_id: &str, + expected_name: &str, + origin: &Url, +) -> bool { + observer.credential_instance_id == credential_instance_id + && observer.name == expected_name + && observer.protocol_version == 2 + && !observer.key.is_empty() + && !observer.prefix.is_empty() + && !observer.name.is_empty() + && !observer.ingest_url.is_empty() + && !contains_invalid_header_value(&observer.key) + && confine_path(origin, &observer.ingest_url).is_ok() +} + +fn write_observer_durably( + config_root: &Path, + observer: &ObserverState, + fault: &dyn DurableWriteFault, +) -> Result<(), PrivateStateError> { + let bytes = serde_json::to_vec(observer).map_err(|_| PrivateStateError::MalformedObserver)?; + atomic_write_bytes_with_fault(&config_root.join(OBSERVER_FILENAME), &bytes, fault).map_err( + |error| { + map_private_file( + error, + PrivateTargetKind::Observer, + PrivateIoOperation::Persist, + ) + }, + ) +} + #[cfg(test)] pub(crate) fn persist_observer( config_root: &Path, observer: &ObserverState, ) -> Result<(), PrivateStateError> { - let bytes = serde_json::to_vec(observer).map_err(|_| PrivateStateError::MalformedObserver)?; - atomic_write_bytes(&config_root.join(OBSERVER_FILENAME), &bytes).map_err(|error| { - map_private_file( - error, - PrivateTargetKind::Observer, - PrivateIoOperation::Persist, - ) - }) + write_observer_durably(config_root, observer, &NoWriteFault) } fn contains_invalid_header_value(value: &str) -> bool { @@ -497,20 +515,17 @@ fn contains_invalid_header_value(value: &str) -> bool { fn persist_and_publish_observer( config_root: &Path, + credential_instance_id: &str, + expected_name: &str, + origin: &Url, observer: &ObserverState, opener: &PrivateLinkOpener, fault: &dyn DurableWriteFault, ) -> Result<(), PrivateStateError> { - let bytes = serde_json::to_vec(observer).map_err(|_| PrivateStateError::MalformedObserver)?; - atomic_write_bytes_with_fault(&config_root.join(OBSERVER_FILENAME), &bytes, fault).map_err( - |error| { - map_private_file( - error, - PrivateTargetKind::Observer, - PrivateIoOperation::Persist, - ) - }, - )?; + if !observer_is_valid(observer, credential_instance_id, expected_name, origin) { + return Err(PrivateStateError::RegistrationInvalid); + } + write_observer_durably(config_root, observer, fault)?; opener.set_registered(observer) } @@ -598,6 +613,8 @@ pub(crate) struct PrivateLinkSession { handle: JournalBridgeHandle, token_persistence: Arc, state_lock: PrivateStateLock, + credential_instance_id: String, + expected_name: String, } impl PrivateLinkSession { @@ -631,6 +648,9 @@ impl PrivateLinkSession { fn publish_observer(&self, observer: &ObserverState) -> Result<(), PrivateStateError> { persist_and_publish_observer( self.state_lock.root(), + &self.credential_instance_id, + &self.expected_name, + &self.origin, observer, &self.opener, &NoWriteFault, @@ -643,7 +663,15 @@ impl PrivateLinkSession { observer: &ObserverState, fault: &dyn DurableWriteFault, ) -> Result<(), PrivateStateError> { - persist_and_publish_observer(self.state_lock.root(), observer, &self.opener, fault) + persist_and_publish_observer( + self.state_lock.root(), + &self.credential_instance_id, + &self.expected_name, + &self.origin, + observer, + &self.opener, + fault, + ) } } @@ -838,6 +866,8 @@ async fn start_private_link_session_inner( handle, token_persistence, state_lock, + credential_instance_id, + expected_name: expected_name.to_owned(), }) } @@ -889,6 +919,22 @@ mod tests { } } + fn assert_registered_auth(request: &crate::private_link_test_peer::PeerRequest, key: &str) { + for (name, expected) in [ + (OBSERVER_HEADER_NAME, key.to_owned()), + ("authorization", format!("Bearer {key}")), + (PROTOCOL_VERSION_HEADER_NAME, "2".to_owned()), + ] { + let values = request + .headers + .iter() + .filter(|(candidate, _)| candidate.eq_ignore_ascii_case(name)) + .map(|(_, value)| value.as_str()) + .collect::>(); + assert_eq!(values, vec![expected.as_str()], "{name}"); + } + } + async fn start_peer_session(peer: &PrivateLinkPeer) -> (tempfile::TempDir, PrivateLinkSession) { let temp = tempfile::tempdir().unwrap(); let session = start_private_link_session(temp.path(), peer.credential(), "stream") @@ -996,6 +1042,105 @@ mod tests { opener_rejection_test!(opener_stays_unregistered_for_fragment, observer("/a#f")); opener_rejection_test!(opener_stays_unregistered_for_backslash, observer("/a\\b")); + async fn assert_publish_rejection(mutate: impl FnOnce(&mut ObserverState)) { + let temp = tempfile::tempdir().unwrap(); + let peer = PrivateLinkPeer::start().await; + peer.enqueue_response(200, b"{}".to_vec()); + let credential = peer.credential(); + let mut state = ObserverState { + credential_instance_id: credential.instance_id.clone(), + ..observer("/ingest") + }; + mutate(&mut state); + let session = start_private_link_session(temp.path(), credential, "stream") + .await + .unwrap(); + assert!(matches!( + publish_observer_registration(&session, &state), + Err(PrivateStateError::RegistrationInvalid) + )); + assert!(!temp.path().join(OBSERVER_FILENAME).exists()); + session + .request(Method::GET, "/still-unregistered") + .unwrap() + .send() + .await + .unwrap(); + let requests = peer.requests(); + assert_eq!(requests.len(), 1); + assert_eq!( + requests[0] + .headers + .iter() + .filter(|(name, _)| { + name.eq_ignore_ascii_case(PROTOCOL_VERSION_HEADER_NAME) + || name.eq_ignore_ascii_case(OBSERVER_HEADER_NAME) + || name.eq_ignore_ascii_case("authorization") + }) + .cloned() + .collect::>(), + vec![(PROTOCOL_VERSION_HEADER_NAME.to_owned(), "2".to_owned())], + ); + session.shutdown().await.unwrap(); + peer.shutdown().await; + } + + macro_rules! publish_rejection_test { + ($name:ident, $field:ident, $value:expr) => { + #[tokio::test] + async fn $name() { + assert_publish_rejection(|state| state.$field = $value.into()).await; + } + }; + } + + publish_rejection_test!( + publish_rejects_credential_instance_mismatch, + credential_instance_id, + "other" + ); + publish_rejection_test!(publish_rejects_name_mismatch, name, "other"); + #[tokio::test] + async fn publish_rejects_unsupported_protocol() { + assert_publish_rejection(|state| state.protocol_version = 3).await; + } + publish_rejection_test!(publish_rejects_empty_key, key, ""); + publish_rejection_test!(publish_rejects_unsafe_key, key, "bad\u{1}key"); + publish_rejection_test!(publish_rejects_empty_prefix, prefix, ""); + publish_rejection_test!(publish_rejects_empty_name, name, ""); + publish_rejection_test!(publish_rejects_empty_ingest_path, ingest_url, ""); + publish_rejection_test!(publish_rejects_relative_ingest_path, ingest_url, "relative"); + publish_rejection_test!( + publish_rejects_scheme_relative_ingest_path, + ingest_url, + "//host/x" + ); + publish_rejection_test!(publish_rejects_raw_ingest_traversal, ingest_url, "/a/../b"); + publish_rejection_test!( + publish_rejects_encoded_ingest_traversal, + ingest_url, + "/a/%2e%2e/b" + ); + publish_rejection_test!( + publish_rejects_mixed_encoded_ingest_traversal, + ingest_url, + "/a/%2E./b" + ); + publish_rejection_test!(publish_rejects_encoded_ingest_slash, ingest_url, "/a/%2f/b"); + publish_rejection_test!( + publish_rejects_encoded_ingest_backslash, + ingest_url, + "/a/%5c/b" + ); + publish_rejection_test!( + publish_rejects_double_encoded_ingest_traversal, + ingest_url, + "/a/%252e%252e/b" + ); + publish_rejection_test!(publish_rejects_query_ingest_path, ingest_url, "/a?q"); + publish_rejection_test!(publish_rejects_fragment_ingest_path, ingest_url, "/a#f"); + publish_rejection_test!(publish_rejects_backslash_ingest_path, ingest_url, "/a\\b"); + async fn assert_redacted_setup_rejection(bytes: &[u8]) { let temp = tempfile::tempdir().unwrap(); let pairer = FakePairer { @@ -1408,18 +1553,7 @@ mod tests { .unwrap(); let requests = peer.requests(); assert_eq!(requests.len(), 1); - for name in [ - OBSERVER_HEADER_NAME, - PROTOCOL_VERSION_HEADER_NAME, - "authorization", - ] { - assert!( - requests[0] - .headers - .iter() - .any(|(candidate, _)| candidate.eq_ignore_ascii_case(name)) - ); - } + assert_registered_auth(&requests[0], "observer-key"); session.shutdown().await.unwrap(); peer.shutdown().await; } @@ -1427,6 +1561,7 @@ mod tests { #[tokio::test] async fn bridge_reuses_one_carrier_across_registration_transition() { let peer = PrivateLinkPeer::start().await; + let credential_instance_id = peer.credential().instance_id; peer.enqueue_response(200, b"{}".to_vec()); peer.enqueue_response(200, b"{}".to_vec()); let (_temp, session) = start_peer_session(&peer).await; @@ -1440,7 +1575,14 @@ mod tests { .status(), StatusCode::OK ); - publish_observer_registration(&session, &observer("/app/observer/ingest")).unwrap(); + publish_observer_registration( + &session, + &ObserverState { + credential_instance_id, + ..observer("/app/observer/ingest") + }, + ) + .unwrap(); assert_eq!( session .request(Method::GET, "/registered") @@ -1453,6 +1595,9 @@ mod tests { ); let requests = peer.requests(); assert_eq!(requests.len(), 2); + assert_eq!(requests[0].method, "GET"); + assert_eq!(requests[0].path, "/unregistered"); + assert!(requests[0].body.is_empty()); assert_eq!( requests[0] .headers @@ -1468,18 +1613,10 @@ mod tests { .any(|(name, _)| name.eq_ignore_ascii_case(OBSERVER_HEADER_NAME) || name.eq_ignore_ascii_case("authorization")) ); - for name in [ - OBSERVER_HEADER_NAME, - PROTOCOL_VERSION_HEADER_NAME, - "authorization", - ] { - assert!( - requests[1] - .headers - .iter() - .any(|(candidate, _)| candidate.eq_ignore_ascii_case(name)) - ); - } + assert_eq!(requests[1].method, "GET"); + assert_eq!(requests[1].path, "/registered"); + assert!(requests[1].body.is_empty()); + assert_registered_auth(&requests[1], "observer-key"); assert_eq!(peer.accepted_carriers(), 1); session.shutdown().await.unwrap(); peer.shutdown().await; @@ -1640,6 +1777,7 @@ mod tests { persist_observer(temp.path(), &prior).unwrap(); let prior_bytes = fs::read(temp.path().join(OBSERVER_FILENAME)).unwrap(); let peer = PrivateLinkPeer::start().await; + let credential_instance_id = peer.credential().instance_id; for _ in 0..6 { peer.enqueue_response(200, b"{}".to_vec()); } @@ -1647,6 +1785,7 @@ mod tests { .await .unwrap(); let next = ObserverState { + credential_instance_id, key: "new-observer-key".into(), ingest_url: "/new".into(), ..observer("/new") @@ -1676,18 +1815,8 @@ mod tests { } else { assert_eq!(current_bytes, prior_bytes); } - let loaded = load_observer( - temp.path(), - "instance", - "stream", - &Url::parse("http://127.0.0.1:1").unwrap(), - ) - .unwrap(); - if stage == DurableWriteStage::DirSync { - assert!(loaded.as_ref() == Some(&prior) || loaded.as_ref() == Some(&next)); - } else { - assert!(loaded.as_ref() == Some(&prior)); - } + let decoded: ObserverState = serde_json::from_slice(¤t_bytes).unwrap(); + assert!(decoded == prior || decoded == next); session .request(Method::GET, "/still-unregistered") .unwrap() @@ -1713,18 +1842,7 @@ mod tests { || name.eq_ignore_ascii_case("authorization") })); } - for name in [ - OBSERVER_HEADER_NAME, - PROTOCOL_VERSION_HEADER_NAME, - "authorization", - ] { - assert!( - requests[5] - .headers - .iter() - .any(|(candidate, _)| candidate.eq_ignore_ascii_case(name)) - ); - } + assert_registered_auth(&requests[5], "new-observer-key"); session.shutdown().await.unwrap(); peer.shutdown().await; } @@ -1826,6 +1944,7 @@ mod tests { let capability = capability.lock().unwrap().clone().unwrap(); let observer_key = "observer-key-sentinel"; let registered = ObserverState { + credential_instance_id: session.credential_instance_id.clone(), key: observer_key.into(), ..observer("/ingest") }; diff --git a/crates/solstone-linux/src/private_link_test_peer.rs b/crates/solstone-linux/src/private_link_test_peer.rs index a2e8e44..7d0da4e 100644 --- a/crates/solstone-linux/src/private_link_test_peer.rs +++ b/crates/solstone-linux/src/private_link_test_peer.rs @@ -36,7 +36,10 @@ use tokio_rustls::{TlsAcceptor, server::TlsStream}; #[derive(Clone)] pub(crate) struct PeerRequest { + pub(crate) method: String, + pub(crate) path: String, pub(crate) headers: Vec<(String, String)>, + pub(crate) body: Vec, } #[derive(Clone)] @@ -300,11 +303,16 @@ fn parse_request(raw: &[u8]) -> Option { let head = std::str::from_utf8(&raw[..split]).ok()?; let mut lines = head.split("\r\n"); let mut request = lines.next()?.split_whitespace(); - request.next()?; - request.next()?; + let method = request.next()?.to_owned(); + let path = request.next()?.to_owned(); let headers = lines .filter_map(|line| line.split_once(':')) .map(|(name, value)| (name.to_owned(), value.trim().to_owned())) .collect(); - Some(PeerRequest { headers }) + Some(PeerRequest { + method, + path, + headers, + body: raw[split + 4..].to_vec(), + }) }