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::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
16
17use crate::error::Result;
18use 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.
22pub 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.
44pub async fn serve<S>(
45 repo_path: &Path,
46 service: Service,
47 protocol_v2: bool,
48 mut stream: S,
49) -> Result<()>
50where
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
59async fn upload_pack<S>(repo_path: &Path, protocol_v2: bool, stream: &mut S) -> Result<()>
60where
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
100async fn receive_pack<S>(repo_path: &Path, stream: S) -> Result<()>
101where
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`).
123fn 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).
133async fn read_until_flush<S>(stream: &mut S) -> Result<Vec<u8>>
134where
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.
157async fn read_until_done<S>(stream: &mut S) -> Result<Vec<u8>>
158where
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
180enum 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.
188async fn read_pkt<S>(stream: &mut S) -> Result<Option<Pkt>>
189where
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
217fn proto_err(message: &'static str) -> crate::error::Error {
218 std::io::Error::new(std::io::ErrorKind::InvalidData, message).into()
219}