anvilsign in

collin/anvil

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
12use std::path::Path;
13
14use gitserver_core::backend::GitBackend;
15use tokio::io::{
16 AsyncRead,
17 AsyncReadExt,
18 AsyncWrite,
19 AsyncWriteExt,
20};
21
22use crate::error::Result;
23use 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.
31pub 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.
53pub async fn serve<S>(
54 repo_path: &Path,
55 service: Service,
56 protocol_v2: bool,
57 mut stream: S,
58) -> Result<()>
59where
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
68async fn upload_pack<S>(repo_path: &Path, protocol_v2: bool, stream: &mut S) -> Result<()>
69where
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
109async fn receive_pack<S>(repo_path: &Path, stream: S) -> Result<()>
110where
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`).
132fn 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).
142async fn read_until_flush<S>(stream: &mut S) -> Result<Vec<u8>>
143where
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.
166async fn read_until_done<S>(stream: &mut S) -> Result<Vec<u8>>
167where
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
189enum 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.
197async fn read_pkt<S>(stream: &mut S) -> Result<Option<Pkt>>
198where
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
226fn proto_err(message: &'static str) -> crate::error::Error {
227 std::io::Error::new(std::io::ErrorKind::InvalidData, message).into()
228}