| 1 | //! Git smart protocol over a bidirectional stream (SSH transport / `git://`). |
| 2 | //! |
| 3 | //! Unlike smart-HTTP — which is a sequence of independent request/response pairs |
| 4 | //! — the SSH transport is one duplex stream: the server writes the ref |
| 5 | //! advertisement, then reads the client's request(s) and writes the |
| 6 | //! response/packfile on the same connection. SSH also omits the HTTP-only |
| 7 | //! `# service=…` announcement line. |
| 8 | //! |
| 9 | //! This module reuses the same [`crate::smart_http`] engine; it only adds the |
| 10 | //! stream framing (reading pkt-line requests and stripping the HTTP prefix). |
| 11 | |
| 12 | use std::path::Path; |
| 13 | |
| 14 | use gitserver_core::backend::GitBackend; |
| 15 | use tokio::io::{ |
| 16 | AsyncRead, |
| 17 | AsyncReadExt, |
| 18 | AsyncWrite, |
| 19 | AsyncWriteExt, |
| 20 | }; |
| 21 | |
| 22 | use crate::error::Result; |
| 23 | use crate::smart_http::{ |
| 24 | self, |
| 25 | Service, |
| 26 | UploadPack, |
| 27 | }; |
| 28 | |
| 29 | /// Parse an SSH `exec` command such as `git-upload-pack '/alice/hello.git'`, |
| 30 | /// returning the service and the (leading-slash-trimmed) repository path. |
| 31 | pub fn parse_command(command: &str) -> Option<(Service, String)> { |
| 32 | let command = command.trim(); |
| 33 | for (prefix, service) in [ |
| 34 | ("git-upload-pack", Service::UploadPack), |
| 35 | ("git-receive-pack", Service::ReceivePack), |
| 36 | ] { |
| 37 | if let Some(rest) = command.strip_prefix(prefix) { |
| 38 | let path = rest |
| 39 | .trim() |
| 40 | .trim_matches('\'') |
| 41 | .trim_matches('"') |
| 42 | .trim_start_matches('/'); |
| 43 | return Some((service, path.to_string())); |
| 44 | } |
| 45 | } |
| 46 | None |
| 47 | } |
| 48 | |
| 49 | /// Serve one git request over `stream` for the bare repo at `repo_path`. |
| 50 | /// |
| 51 | /// Takes ownership of the stream: `receive-pack` splits it so the protocol |
| 52 | /// engine can own the read half. |
| 53 | pub async fn serve<S>( |
| 54 | repo_path: &Path, |
| 55 | service: Service, |
| 56 | protocol_v2: bool, |
| 57 | mut stream: S, |
| 58 | ) -> Result<()> |
| 59 | where |
| 60 | S: AsyncRead + AsyncWrite + Unpin + Send + 'static, |
| 61 | { |
| 62 | match service { |
| 63 | Service::UploadPack => upload_pack(repo_path, protocol_v2, &mut stream).await, |
| 64 | Service::ReceivePack => receive_pack(repo_path, stream).await, |
| 65 | } |
| 66 | } |
| 67 | |
| 68 | async fn upload_pack<S>(repo_path: &Path, protocol_v2: bool, stream: &mut S) -> Result<()> |
| 69 | where |
| 70 | S: AsyncRead + AsyncWrite + Unpin + Send, |
| 71 | { |
| 72 | let advertisement = smart_http::advertise(repo_path, Service::UploadPack, protocol_v2)?; |
| 73 | stream |
| 74 | .write_all(strip_http_service_prefix(&advertisement)) |
| 75 | .await?; |
| 76 | stream.flush().await?; |
| 77 | |
| 78 | if protocol_v2 { |
| 79 | // v2: a sequence of commands (ls-refs, fetch), each terminated by flush. |
| 80 | loop { |
| 81 | let command = read_until_flush(stream).await?; |
| 82 | if command.is_empty() { |
| 83 | break; // client closed the connection |
| 84 | } |
| 85 | match smart_http::upload_pack_v2(repo_path, &command).await? { |
| 86 | UploadPack::Buffered(body) => { |
| 87 | stream.write_all(&body).await?; |
| 88 | stream.flush().await?; |
| 89 | } |
| 90 | UploadPack::Pack(mut pack) => { |
| 91 | tokio::io::copy(&mut pack, stream).await?; |
| 92 | stream.flush().await?; |
| 93 | break; |
| 94 | } |
| 95 | } |
| 96 | } |
| 97 | } else { |
| 98 | // v0/v1: a single want/have request terminated by `done`. |
| 99 | let request = read_until_done(stream).await?; |
| 100 | if !request.is_empty() { |
| 101 | let mut pack = smart_http::upload_pack_v0(repo_path, &request).await?; |
| 102 | tokio::io::copy(&mut pack, stream).await?; |
| 103 | stream.flush().await?; |
| 104 | } |
| 105 | } |
| 106 | Ok(()) |
| 107 | } |
| 108 | |
| 109 | async fn receive_pack<S>(repo_path: &Path, stream: S) -> Result<()> |
| 110 | where |
| 111 | S: AsyncRead + AsyncWrite + Unpin + Send + 'static, |
| 112 | { |
| 113 | let (read_half, mut write_half) = tokio::io::split(stream); |
| 114 | let backend = GitBackend::new(repo_path.to_path_buf()); |
| 115 | |
| 116 | // Advertise refs (no HTTP `# service` prefix on the SSH transport). |
| 117 | let advertisement = backend.advertise_receive_refs()?; |
| 118 | write_half.write_all(&advertisement).await?; |
| 119 | write_half.flush().await?; |
| 120 | |
| 121 | // The engine reads exactly the commands + packfile (not to EOF), so this |
| 122 | // does not deadlock waiting for the client to half-close. |
| 123 | let report = backend.receive_pack(read_half).await?; |
| 124 | write_half.write_all(&report).await?; |
| 125 | write_half.flush().await?; |
| 126 | Ok(()) |
| 127 | } |
| 128 | |
| 129 | /// The HTTP advertisements begin with a `# service=…` pkt-line followed by a |
| 130 | /// flush; the SSH/`git://` transports omit it. Return the bytes after the first |
| 131 | /// flush packet (the service comment never contains `0000`). |
| 132 | fn strip_http_service_prefix(advertisement: &[u8]) -> &[u8] { |
| 133 | match advertisement.windows(4).position(|w| w == b"0000") { |
| 134 | Some(pos) => &advertisement[pos + 4..], |
| 135 | None => advertisement, |
| 136 | } |
| 137 | } |
| 138 | |
| 139 | /// Read raw pkt-lines until (and including) a flush packet. Returns the exact |
| 140 | /// bytes read, suitable for the v2 command parser. An empty result means EOF |
| 141 | /// before any data (the client closed the connection). |
| 142 | async fn read_until_flush<S>(stream: &mut S) -> Result<Vec<u8>> |
| 143 | where |
| 144 | S: AsyncRead + Unpin, |
| 145 | { |
| 146 | let mut buf = Vec::new(); |
| 147 | loop { |
| 148 | match read_pkt(stream).await? { |
| 149 | None => break, |
| 150 | Some(Pkt::Flush) => { |
| 151 | buf.extend_from_slice(b"0000"); |
| 152 | break; |
| 153 | } |
| 154 | Some(Pkt::Delim) => buf.extend_from_slice(b"0001"), |
| 155 | Some(Pkt::ResponseEnd) => buf.extend_from_slice(b"0002"), |
| 156 | Some(Pkt::Line(len, payload)) => { |
| 157 | buf.extend_from_slice(format!("{len:04x}").as_bytes()); |
| 158 | buf.extend_from_slice(&payload); |
| 159 | } |
| 160 | } |
| 161 | } |
| 162 | Ok(buf) |
| 163 | } |
| 164 | |
| 165 | /// Read raw pkt-lines until a `done` line (v0/v1 negotiation end) or EOF. |
| 166 | async fn read_until_done<S>(stream: &mut S) -> Result<Vec<u8>> |
| 167 | where |
| 168 | S: AsyncRead + Unpin, |
| 169 | { |
| 170 | let mut buf = Vec::new(); |
| 171 | loop { |
| 172 | match read_pkt(stream).await? { |
| 173 | None => break, |
| 174 | Some(Pkt::Flush) => buf.extend_from_slice(b"0000"), |
| 175 | Some(Pkt::Delim) => buf.extend_from_slice(b"0001"), |
| 176 | Some(Pkt::ResponseEnd) => buf.extend_from_slice(b"0002"), |
| 177 | Some(Pkt::Line(len, payload)) => { |
| 178 | buf.extend_from_slice(format!("{len:04x}").as_bytes()); |
| 179 | buf.extend_from_slice(&payload); |
| 180 | if payload.trim_ascii_end() == b"done" { |
| 181 | break; |
| 182 | } |
| 183 | } |
| 184 | } |
| 185 | } |
| 186 | Ok(buf) |
| 187 | } |
| 188 | |
| 189 | enum Pkt { |
| 190 | Flush, |
| 191 | Delim, |
| 192 | ResponseEnd, |
| 193 | Line(usize, Vec<u8>), |
| 194 | } |
| 195 | |
| 196 | /// Read a single pkt-line. Returns `None` on a clean EOF at a packet boundary. |
| 197 | async fn read_pkt<S>(stream: &mut S) -> Result<Option<Pkt>> |
| 198 | where |
| 199 | S: AsyncRead + Unpin, |
| 200 | { |
| 201 | let mut header = [0u8; 4]; |
| 202 | match stream.read_exact(&mut header).await { |
| 203 | Ok(_) => {} |
| 204 | Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => return Ok(None), |
| 205 | Err(e) => return Err(e.into()), |
| 206 | } |
| 207 | let len = usize::from_str_radix( |
| 208 | std::str::from_utf8(&header).map_err(|_| proto_err("non-utf8 pkt-line length"))?, |
| 209 | 16, |
| 210 | ) |
| 211 | .map_err(|_| proto_err("invalid pkt-line length"))?; |
| 212 | |
| 213 | match len { |
| 214 | 0 => Ok(Some(Pkt::Flush)), |
| 215 | 1 => Ok(Some(Pkt::Delim)), |
| 216 | 2 => Ok(Some(Pkt::ResponseEnd)), |
| 217 | 3 => Err(proto_err("invalid pkt-line length 0003")), |
| 218 | _ => { |
| 219 | let mut payload = vec![0u8; len - 4]; |
| 220 | stream.read_exact(&mut payload).await?; |
| 221 | Ok(Some(Pkt::Line(len, payload))) |
| 222 | } |
| 223 | } |
| 224 | } |
| 225 | |
| 226 | fn proto_err(message: &'static str) -> crate::error::Error { |
| 227 | std::io::Error::new(std::io::ErrorKind::InvalidData, message).into() |
| 228 | } |