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