| 1 | use std::path::PathBuf; |
| 2 | use std::pin::Pin; |
| 3 | use std::sync::Arc; |
| 4 | use std::sync::atomic::{AtomicBool, Ordering}; |
| 5 | use std::task::{Context, Poll}; |
| 6 | |
| 7 | use tokio::io::AsyncRead; |
| 8 | use tokio::time::{Duration, Sleep, sleep}; |
| 9 | use tokio_util::io::SyncIoBridge; |
| 10 | |
| 11 | use crate::error::Result; |
| 12 | use crate::pack::UploadPackRequest; |
| 13 | |
| 14 | pub const RECEIVE_PACK_TIMEOUT: Duration = Duration::from_secs(300); |
| 15 | const RECEIVE_PACK_IDLE_TIMEOUT: Duration = Duration::from_secs(30); |
| 16 | |
| 17 | struct TimedAsyncRead<R> { |
| 18 | inner: R, |
| 19 | timeout: Duration, |
| 20 | sleep: Option<Pin<Box<Sleep>>>, |
| 21 | interrupt: Arc<AtomicBool>, |
| 22 | } |
| 23 | |
| 24 | impl<R> TimedAsyncRead<R> { |
| 25 | fn new(inner: R, timeout: Duration, interrupt: Arc<AtomicBool>) -> Self { |
| 26 | Self { |
| 27 | inner, |
| 28 | timeout, |
| 29 | sleep: None, |
| 30 | interrupt, |
| 31 | } |
| 32 | } |
| 33 | } |
| 34 | |
| 35 | impl<R> AsyncRead for TimedAsyncRead<R> |
| 36 | where |
| 37 | R: AsyncRead + Unpin, |
| 38 | { |
| 39 | fn poll_read( |
| 40 | mut self: Pin<&mut Self>, |
| 41 | cx: &mut Context<'_>, |
| 42 | buf: &mut tokio::io::ReadBuf<'_>, |
| 43 | ) -> Poll<std::io::Result<()>> { |
| 44 | if self.sleep.is_none() { |
| 45 | self.sleep = Some(Box::pin(sleep(self.timeout))); |
| 46 | } |
| 47 | |
| 48 | let before = buf.filled().len(); |
| 49 | match Pin::new(&mut self.inner).poll_read(cx, buf) { |
| 50 | Poll::Ready(Ok(())) => { |
| 51 | if buf.filled().len() > before { |
| 52 | self.sleep = None; |
| 53 | } |
| 54 | Poll::Ready(Ok(())) |
| 55 | } |
| 56 | Poll::Ready(Err(err)) => Poll::Ready(Err(err)), |
| 57 | Poll::Pending => { |
| 58 | if self |
| 59 | .sleep |
| 60 | .as_mut() |
| 61 | .expect("timeout sleep must exist") |
| 62 | .as_mut() |
| 63 | .poll(cx) |
| 64 | .is_ready() |
| 65 | { |
| 66 | self.interrupt.store(true, Ordering::Relaxed); |
| 67 | Poll::Ready(Err(std::io::Error::new( |
| 68 | std::io::ErrorKind::TimedOut, |
| 69 | "receive-pack read timed out", |
| 70 | ))) |
| 71 | } else { |
| 72 | Poll::Pending |
| 73 | } |
| 74 | } |
| 75 | } |
| 76 | } |
| 77 | } |
| 78 | |
| 79 | pub struct GitBackend { |
| 80 | repo_path: PathBuf, |
| 81 | } |
| 82 | |
| 83 | impl GitBackend { |
| 84 | pub fn new(repo_path: PathBuf) -> Self { |
| 85 | Self { repo_path } |
| 86 | } |
| 87 | |
| 88 | pub fn advertise_refs(&self) -> Result<Vec<u8>> { |
| 89 | crate::refs::advertise_refs(&self.repo_path) |
| 90 | } |
| 91 | |
| 92 | pub fn advertise_receive_refs(&self) -> Result<Vec<u8>> { |
| 93 | crate::receive_pack::advertise_receive_refs(&self.repo_path) |
| 94 | } |
| 95 | |
| 96 | pub async fn upload_pack(&self, request: &UploadPackRequest) -> Result<impl AsyncRead + use<>> { |
| 97 | crate::pack::generate_pack(&self.repo_path, request) |
| 98 | } |
| 99 | |
| 100 | pub async fn receive_pack<R>(&self, request: R) -> Result<Vec<u8>> |
| 101 | where |
| 102 | R: AsyncRead + Unpin + Send + 'static, |
| 103 | { |
| 104 | self.receive_pack_with_timeout(request, RECEIVE_PACK_TIMEOUT) |
| 105 | .await |
| 106 | } |
| 107 | |
| 108 | async fn receive_pack_with_timeout<R>( |
| 109 | &self, |
| 110 | request: R, |
| 111 | timeout_duration: Duration, |
| 112 | ) -> Result<Vec<u8>> |
| 113 | where |
| 114 | R: AsyncRead + Unpin + Send + 'static, |
| 115 | { |
| 116 | let repo_path = self.repo_path.clone(); |
| 117 | let interrupt = Arc::new(AtomicBool::new(false)); |
| 118 | let watchdog_interrupt = interrupt.clone(); |
| 119 | let watchdog = tokio::spawn(async move { |
| 120 | sleep(timeout_duration).await; |
| 121 | watchdog_interrupt.store(true, Ordering::Relaxed); |
| 122 | }); |
| 123 | |
| 124 | let join = tokio::task::spawn_blocking(move || { |
| 125 | let request = |
| 126 | TimedAsyncRead::new(request, RECEIVE_PACK_IDLE_TIMEOUT, interrupt.clone()); |
| 127 | let mut request = SyncIoBridge::new(request); |
| 128 | crate::receive_pack::receive_pack_with_interrupt( |
| 129 | &repo_path, |
| 130 | &mut request, |
| 131 | interrupt.as_ref(), |
| 132 | ) |
| 133 | }) |
| 134 | .await |
| 135 | .map_err(|e| crate::error::Error::Protocol(format!("receive-pack task panicked: {e}"))); |
| 136 | |
| 137 | watchdog.abort(); |
| 138 | join? |
| 139 | } |
| 140 | } |
| 141 | |
| 142 | #[cfg(test)] |
| 143 | mod tests { |
| 144 | use std::process::Command; |
| 145 | |
| 146 | use tempfile::TempDir; |
| 147 | |
| 148 | use super::*; |
| 149 | |
| 150 | fn create_repo_with_commit(root: &std::path::Path) -> PathBuf { |
| 151 | let repo_path = root.join("test.git"); |
| 152 | let work_dir = root.join("work"); |
| 153 | std::fs::create_dir(&work_dir).unwrap(); |
| 154 | Command::new("git") |
| 155 | .args(["init", "--bare", repo_path.to_str().unwrap()]) |
| 156 | .output() |
| 157 | .unwrap(); |
| 158 | Command::new("git") |
| 159 | .args(["symbolic-ref", "HEAD", "refs/heads/main"]) |
| 160 | .current_dir(&repo_path) |
| 161 | .output() |
| 162 | .unwrap(); |
| 163 | Command::new("git") |
| 164 | .args([ |
| 165 | "clone", |
| 166 | repo_path.to_str().unwrap(), |
| 167 | work_dir.to_str().unwrap(), |
| 168 | ]) |
| 169 | .output() |
| 170 | .unwrap(); |
| 171 | Command::new("git") |
| 172 | .current_dir(&work_dir) |
| 173 | .args(["commit", "--allow-empty", "-m", "init"]) |
| 174 | .env("GIT_AUTHOR_NAME", "Test") |
| 175 | .env("GIT_AUTHOR_EMAIL", "t@t.com") |
| 176 | .env("GIT_COMMITTER_NAME", "Test") |
| 177 | .env("GIT_COMMITTER_EMAIL", "t@t.com") |
| 178 | .output() |
| 179 | .unwrap(); |
| 180 | Command::new("git") |
| 181 | .current_dir(&work_dir) |
| 182 | .args(["push", "origin", "main"]) |
| 183 | .output() |
| 184 | .unwrap(); |
| 185 | repo_path |
| 186 | } |
| 187 | |
| 188 | #[test] |
| 189 | fn backend_advertise_refs() { |
| 190 | let root = TempDir::new().unwrap(); |
| 191 | let repo_path = create_repo_with_commit(root.path()); |
| 192 | let backend = GitBackend::new(repo_path); |
| 193 | let output = backend.advertise_refs().unwrap(); |
| 194 | let output_str = String::from_utf8_lossy(&output); |
| 195 | assert!(output_str.contains("refs/heads/main")); |
| 196 | } |
| 197 | |
| 198 | #[tokio::test] |
| 199 | async fn backend_upload_pack() { |
| 200 | let root = TempDir::new().unwrap(); |
| 201 | let repo_path = create_repo_with_commit(root.path()); |
| 202 | let repo = gix::open(&repo_path).unwrap(); |
| 203 | let head = repo.head_id().unwrap(); |
| 204 | |
| 205 | let backend = GitBackend::new(repo_path); |
| 206 | let request = UploadPackRequest { |
| 207 | wants: vec![head.detach()], |
| 208 | haves: vec![], |
| 209 | done: true, |
| 210 | capabilities: Default::default(), |
| 211 | shallow: Default::default(), |
| 212 | object_ids: None, |
| 213 | }; |
| 214 | let reader = backend.upload_pack(&request).await.unwrap(); |
| 215 | let mut buf = Vec::new(); |
| 216 | tokio::io::AsyncReadExt::read_to_end(&mut tokio::io::BufReader::new(reader), &mut buf) |
| 217 | .await |
| 218 | .unwrap(); |
| 219 | assert!(buf.windows(4).any(|w| w == b"PACK")); |
| 220 | } |
| 221 | |
| 222 | #[tokio::test] |
| 223 | async fn backend_receive_pack_times_out_on_stalled_reader() { |
| 224 | let root = TempDir::new().unwrap(); |
| 225 | let repo_path = create_repo_with_commit(root.path()); |
| 226 | let backend = GitBackend::new(repo_path); |
| 227 | let (reader, _writer) = tokio::io::duplex(1); |
| 228 | |
| 229 | let err = backend |
| 230 | .receive_pack_with_timeout(reader, Duration::from_millis(50)) |
| 231 | .await |
| 232 | .unwrap_err(); |
| 233 | |
| 234 | match err { |
| 235 | crate::error::Error::Io(inner) => { |
| 236 | assert_eq!(inner.kind(), std::io::ErrorKind::TimedOut); |
| 237 | } |
| 238 | other => panic!("expected timeout io error, got {other}"), |
| 239 | } |
| 240 | } |
| 241 | } |