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