Skip to content
286 linesCodeBlameRaw

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Initial g1t: services, event bus, intents and attempts1//! SSH front end for g1t. Accepts `git@g1t.sh:owner/repo.git`, authenticates
2//! the client's public key against the g1t Worker, and bridges git to
3//! Artifacts over HTTPS.
4//!
5//! Connections arrive either as raw TCP or wrapped in a WebSocket, which is
6//! how they reach a Cloudflare Container when tunnelled through a Worker.
Artifacts: an hourly view of operations and errors, and SSH must meter its git before it ships7//!
8//! Not deployed. Before it ships, its git operations must be metered: the
9//! bridge reaches the store directly, so `git_http` never counts them (see
10//! docs/ARTIFACTS.md, "where the gap came from").
Initial g1t: services, event bus, intents and attempts11
12mod api;
13mod git;
14
15use std::collections::HashMap;
16use std::sync::Arc;
17use std::time::Duration;
18
19use anyhow::{Context, Result};
20use futures_util::{SinkExt, StreamExt};
API and MCP server, Rust identity service, registration, site redesign21use russh::keys::PrivateKey;
Initial g1t: services, event bus, intents and attempts22use russh::keys::ssh_key::{HashAlg, PublicKey};
23use russh::server::{Auth, Config, Handler, Msg, Session};
24use russh::{Channel, ChannelId, MethodKind, MethodSet};
25use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt, BufReader};
26use tokio::net::TcpListener;
27use tokio_tungstenite::tungstenite::Message;
28
29use api::{Api, Service, User};
30
31const TCP_ADDR: &str = "0.0.0.0:2222";
32const WEBSOCKET_ADDR: &str = "0.0.0.0:8080";
33
34#[tokio::main]
35async fn main() -> Result<()> {
36 let host_key = std::env::var("SSH_HOST_KEY").context("SSH_HOST_KEY is not set")?;
37 let api = Arc::new(Api::new(
38 std::env::var("G1T_API").unwrap_or_else(|_| "https://g1t.sh".into()),
39 std::env::var("INTERNAL_SECRET").context("INTERNAL_SECRET is not set")?,
40 ));
41 let config = Arc::new(Config {
42 keys: vec![PrivateKey::from_openssh(host_key).context("invalid SSH_HOST_KEY")?],
43 methods: MethodSet::from(&[MethodKind::PublicKey][..]),
44 auth_rejection_time: Duration::from_secs(1),
45 auth_rejection_time_initial: Some(Duration::ZERO),
46 inactivity_timeout: Some(Duration::from_secs(300)),
47 ..Default::default()
48 });
49
50 let tcp = TcpListener::bind(TCP_ADDR).await?;
51 let websocket = TcpListener::bind(WEBSOCKET_ADDR).await?;
52 eprintln!("g1t-sshd listening on {TCP_ADDR} (tcp) and {WEBSOCKET_ADDR} (websocket)");
53 loop {
54 tokio::select! {
55 accepted = tcp.accept() => {
56 let (stream, _) = accepted?;
57 tokio::spawn(run_session(config.clone(), api.clone(), stream));
58 }
59 accepted = websocket.accept() => {
60 let (stream, _) = accepted?;
61 let (config, api) = (config.clone(), api.clone());
62 tokio::spawn(async move {
63 match tokio_tungstenite::accept_async(stream).await {
64 Ok(socket) => run_session(config, api, websocket_stream(socket)).await,
65 Err(error) => eprintln!("websocket handshake failed: {error}"),
66 }
67 });
68 }
69 }
70 }
71}
72
73async fn run_session<S>(config: Arc<Config>, api: Arc<Api>, stream: S)
74where
75 S: AsyncRead + AsyncWrite + Unpin + Send + 'static,
76{
77 let handler = Connection {
78 api,
79 user: None,
80 git_protocol: None,
81 channels: HashMap::new(),
82 };
83 let result = match russh::server::run_stream(config, stream, handler).await {
84 Ok(session) => session.await,
85 Err(error) => Err(error),
86 };
87 if let Err(error) = result {
88 eprintln!("session ended: {error:#}");
89 }
90}
91
92/// Presents a WebSocket's binary messages as a plain byte stream.
93fn websocket_stream<S>(socket: tokio_tungstenite::WebSocketStream<S>) -> tokio::io::DuplexStream
94where
95 S: AsyncRead + AsyncWrite + Unpin + Send + 'static,
96{
97 let (ours, theirs) = tokio::io::duplex(64 * 1024);
98 let (mut sink, mut source) = socket.split();
99 let (mut from_ssh, mut to_ssh) = tokio::io::split(theirs);
100 tokio::spawn(async move {
101 while let Some(Ok(message)) = source.next().await {
Polish: phones, copy boxes, the plan page, the landing page, a real glide102 if let Message::Binary(bytes) = message
103 && to_ssh.write_all(&bytes).await.is_err() {
Initial g1t: services, event bus, intents and attempts104 break;
105 }
106 }
107 let _ = to_ssh.shutdown().await;
108 });
109 tokio::spawn(async move {
110 let mut buffer = vec![0u8; 32 * 1024];
111 loop {
112 match from_ssh.read(&mut buffer).await {
113 Ok(0) | Err(_) => break,
114 Ok(read) => {
115 let message = Message::Binary(buffer[..read].to_vec().into());
116 if sink.send(message).await.is_err() {
117 break;
118 }
119 }
120 }
121 }
122 let _ = sink.close().await;
123 });
124 ours
125}
126
127struct Connection {
128 api: Arc<Api>,
129 user: Option<User>,
130 /// Value of the `GIT_PROTOCOL` environment variable, if the client sent it.
131 git_protocol: Option<String>,
132 channels: HashMap<ChannelId, Channel<Msg>>,
133}
134
135impl Connection {
Deploy keys: SSH keys that reach one repository136 /// Who the key signs in as. `used` once the client has proved it holds
137 /// the private key: only then is the key's last use recorded.
138 async fn lookup(&self, key: &PublicKey, used: bool) -> Result<Option<User>> {
Initial g1t: services, event bus, intents and attempts139 let fingerprint = key.fingerprint(HashAlg::Sha256).to_string();
Deploy keys: SSH keys that reach one repository140 self.api.user_for_key(&fingerprint, used).await
Initial g1t: services, event bus, intents and attempts141 }
142}
143
144impl Handler for Connection {
145 type Error = anyhow::Error;
146
147 async fn auth_publickey_offered(&mut self, _: &str, key: &PublicKey) -> Result<Auth> {
Deploy keys: SSH keys that reach one repository148 Ok(match self.lookup(key, false).await? {
Initial g1t: services, event bus, intents and attempts149 Some(_) => Auth::Accept,
150 None => Auth::reject(),
151 })
152 }
153
154 /// Called once the client has proven it holds the private key.
155 async fn auth_publickey(&mut self, _: &str, key: &PublicKey) -> Result<Auth> {
Deploy keys: SSH keys that reach one repository156 self.user = self.lookup(key, true).await?;
Initial g1t: services, event bus, intents and attempts157 Ok(match self.user {
158 Some(_) => Auth::Accept,
159 None => Auth::reject(),
160 })
161 }
162
163 async fn channel_open_session(
164 &mut self,
165 channel: Channel<Msg>,
166 reply: russh::server::ChannelOpenHandle,
167 _: &mut Session,
168 ) -> Result<()> {
169 self.channels.insert(channel.id(), channel);
170 reply.accept().await;
171 Ok(())
172 }
173
174 async fn env_request(
175 &mut self,
176 channel: ChannelId,
177 name: &str,
178 value: &str,
179 session: &mut Session,
180 ) -> Result<()> {
181 if name == "GIT_PROTOCOL" {
182 self.git_protocol = Some(value.to_owned());
183 }
184 session.channel_success(channel)?;
185 Ok(())
186 }
187
188 async fn shell_request(&mut self, id: ChannelId, session: &mut Session) -> Result<()> {
189 let (Some(user), Some(channel)) = (&self.user, self.channels.remove(&id)) else {
190 return Ok(session.channel_failure(id)?);
191 };
192 session.channel_success(id)?;
193 let greeting = format!(
194 "Hi {}! You've successfully authenticated, but g1t does not provide shell access.\r\n",
Deploy keys: SSH keys that reach one repository195 user.greeting_name()
Initial g1t: services, event bus, intents and attempts196 );
197 tokio::spawn(async move {
198 let _ = channel.data(greeting.as_bytes()).await;
199 finish(channel, 1).await;
200 });
201 Ok(())
202 }
203
204 async fn exec_request(
205 &mut self,
206 id: ChannelId,
207 command: &[u8],
208 session: &mut Session,
209 ) -> Result<()> {
210 let (Some(user), Some(channel)) = (self.user.clone(), self.channels.remove(&id)) else {
211 return Ok(session.channel_failure(id)?);
212 };
213 session.channel_success(id)?;
214
215 let command = String::from_utf8_lossy(command).into_owned();
216 let protocol_v2 = self
217 .git_protocol
218 .as_deref()
219 .is_some_and(|value| value.split(':').any(|part| part == "version=2"));
220 let api = self.api.clone();
221 tokio::spawn(async move {
222 let (mut read_half, write_half) = channel.split();
223 let mut reader = BufReader::new(read_half.make_reader());
224 let mut writer = write_half.make_writer();
API and MCP server, Rust identity service, registration, site redesign225 let result =
226 run_git(&api, &user, &command, protocol_v2, &mut reader, &mut writer).await;
Initial g1t: services, event bus, intents and attempts227 let status = match result {
228 Ok(()) => 0,
229 Err(error) => {
230 eprintln!("{}: `{command}` failed: {error:#}", user.username);
231 let _ = writer.write_all(&git::error_pkt("internal error")).await;
232 1
233 }
234 };
235 let _ = writer.shutdown().await;
236 let _ = write_half.exit_status(status).await;
237 let _ = write_half.eof().await;
238 let _ = write_half.close().await;
239 });
240 Ok(())
241 }
242}
243
244async fn finish(channel: Channel<Msg>, status: u32) {
245 let _ = channel.exit_status(status).await;
246 let _ = channel.eof().await;
247 let _ = channel.close().await;
248}
249
250/// Parses `git-upload-pack 'owner/repo.git'` and friends.
251fn parse_command(command: &str) -> Option<(Service, &str, &str)> {
252 let (program, path) = command.trim().rsplit_once(' ')?;
253 let service = match program {
254 "git-upload-pack" | "git upload-pack" => Service::UploadPack,
255 "git-receive-pack" | "git receive-pack" => Service::ReceivePack,
256 _ => return None,
257 };
258 let path = path.trim_matches(['\'', '"']).trim_start_matches('/');
259 let path = path.strip_suffix(".git").unwrap_or(path);
260 let (owner, repo) = path.split_once('/')?;
261 (!owner.is_empty() && !repo.is_empty() && !repo.contains('/')).then_some((service, owner, repo))
262}
263
264async fn run_git<R, W>(
265 api: &Api,
266 user: &User,
267 command: &str,
268 protocol_v2: bool,
269 reader: &mut R,
270 writer: &mut W,
271) -> Result<()>
272where
273 R: tokio::io::AsyncBufRead + Unpin,
274 W: AsyncWrite + Unpin,
275{
276 let Some((service, owner, repo)) = parse_command(command) else {
277 writer
278 .write_all(&git::error_pkt("g1t only supports git over SSH"))
279 .await?;
280 return Ok(());
281 };
282 match api.access(user, owner, repo, service).await? {
283 Ok(access) => git::serve(&api.http, &access, service, protocol_v2, reader, writer).await,
284 Err(message) => Ok(writer.write_all(&git::error_pkt(&message)).await?),
285 }
286}