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 attempts | 1 | //! 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 ships | 7 | //! |
| 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 attempts | 11 | |
| 12 | mod api; | |
| 13 | mod git; | |
| 14 | ||
| 15 | use std::collections::HashMap; | |
| 16 | use std::sync::Arc; | |
| 17 | use std::time::Duration; | |
| 18 | ||
| 19 | use anyhow::{Context, Result}; | |
| 20 | use futures_util::{SinkExt, StreamExt}; | |
| API and MCP server, Rust identity service, registration, site redesign | 21 | use russh::keys::PrivateKey; |
| Initial g1t: services, event bus, intents and attempts | 22 | use russh::keys::ssh_key::{HashAlg, PublicKey}; |
| 23 | use russh::server::{Auth, Config, Handler, Msg, Session}; | |
| 24 | use russh::{Channel, ChannelId, MethodKind, MethodSet}; | |
| 25 | use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt, BufReader}; | |
| 26 | use tokio::net::TcpListener; | |
| 27 | use tokio_tungstenite::tungstenite::Message; | |
| 28 | ||
| 29 | use api::{Api, Service, User}; | |
| 30 | ||
| 31 | const TCP_ADDR: &str = "0.0.0.0:2222"; | |
| 32 | const WEBSOCKET_ADDR: &str = "0.0.0.0:8080"; | |
| 33 | ||
| 34 | #[tokio::main] | |
| 35 | async 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 | ||
| 73 | async fn run_session<S>(config: Arc<Config>, api: Arc<Api>, stream: S) | |
| 74 | where | |
| 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. | |
| 93 | fn websocket_stream<S>(socket: tokio_tungstenite::WebSocketStream<S>) -> tokio::io::DuplexStream | |
| 94 | where | |
| 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 glide | 102 | if let Message::Binary(bytes) = message |
| 103 | && to_ssh.write_all(&bytes).await.is_err() { | |
| Initial g1t: services, event bus, intents and attempts | 104 | 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 | ||
| 127 | struct 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 | ||
| 135 | impl Connection { | |
| 136 | async fn lookup(&self, key: &PublicKey) -> Result<Option<User>> { | |
| 137 | let fingerprint = key.fingerprint(HashAlg::Sha256).to_string(); | |
| 138 | self.api.user_for_key(&fingerprint).await | |
| 139 | } | |
| 140 | } | |
| 141 | ||
| 142 | impl Handler for Connection { | |
| 143 | type Error = anyhow::Error; | |
| 144 | ||
| 145 | async fn auth_publickey_offered(&mut self, _: &str, key: &PublicKey) -> Result<Auth> { | |
| 146 | Ok(match self.lookup(key).await? { | |
| 147 | Some(_) => Auth::Accept, | |
| 148 | None => Auth::reject(), | |
| 149 | }) | |
| 150 | } | |
| 151 | ||
| 152 | /// Called once the client has proven it holds the private key. | |
| 153 | async fn auth_publickey(&mut self, _: &str, key: &PublicKey) -> Result<Auth> { | |
| 154 | self.user = self.lookup(key).await?; | |
| 155 | Ok(match self.user { | |
| 156 | Some(_) => Auth::Accept, | |
| 157 | None => Auth::reject(), | |
| 158 | }) | |
| 159 | } | |
| 160 | ||
| 161 | async fn channel_open_session( | |
| 162 | &mut self, | |
| 163 | channel: Channel<Msg>, | |
| 164 | reply: russh::server::ChannelOpenHandle, | |
| 165 | _: &mut Session, | |
| 166 | ) -> Result<()> { | |
| 167 | self.channels.insert(channel.id(), channel); | |
| 168 | reply.accept().await; | |
| 169 | Ok(()) | |
| 170 | } | |
| 171 | ||
| 172 | async fn env_request( | |
| 173 | &mut self, | |
| 174 | channel: ChannelId, | |
| 175 | name: &str, | |
| 176 | value: &str, | |
| 177 | session: &mut Session, | |
| 178 | ) -> Result<()> { | |
| 179 | if name == "GIT_PROTOCOL" { | |
| 180 | self.git_protocol = Some(value.to_owned()); | |
| 181 | } | |
| 182 | session.channel_success(channel)?; | |
| 183 | Ok(()) | |
| 184 | } | |
| 185 | ||
| 186 | async fn shell_request(&mut self, id: ChannelId, session: &mut Session) -> Result<()> { | |
| 187 | let (Some(user), Some(channel)) = (&self.user, self.channels.remove(&id)) else { | |
| 188 | return Ok(session.channel_failure(id)?); | |
| 189 | }; | |
| 190 | session.channel_success(id)?; | |
| 191 | let greeting = format!( | |
| 192 | "Hi {}! You've successfully authenticated, but g1t does not provide shell access.\r\n", | |
| 193 | user.username | |
| 194 | ); | |
| 195 | tokio::spawn(async move { | |
| 196 | let _ = channel.data(greeting.as_bytes()).await; | |
| 197 | finish(channel, 1).await; | |
| 198 | }); | |
| 199 | Ok(()) | |
| 200 | } | |
| 201 | ||
| 202 | async fn exec_request( | |
| 203 | &mut self, | |
| 204 | id: ChannelId, | |
| 205 | command: &[u8], | |
| 206 | session: &mut Session, | |
| 207 | ) -> Result<()> { | |
| 208 | let (Some(user), Some(channel)) = (self.user.clone(), self.channels.remove(&id)) else { | |
| 209 | return Ok(session.channel_failure(id)?); | |
| 210 | }; | |
| 211 | session.channel_success(id)?; | |
| 212 | ||
| 213 | let command = String::from_utf8_lossy(command).into_owned(); | |
| 214 | let protocol_v2 = self | |
| 215 | .git_protocol | |
| 216 | .as_deref() | |
| 217 | .is_some_and(|value| value.split(':').any(|part| part == "version=2")); | |
| 218 | let api = self.api.clone(); | |
| 219 | tokio::spawn(async move { | |
| 220 | let (mut read_half, write_half) = channel.split(); | |
| 221 | let mut reader = BufReader::new(read_half.make_reader()); | |
| 222 | let mut writer = write_half.make_writer(); | |
| API and MCP server, Rust identity service, registration, site redesign | 223 | let result = |
| 224 | run_git(&api, &user, &command, protocol_v2, &mut reader, &mut writer).await; | |
| Initial g1t: services, event bus, intents and attempts | 225 | let status = match result { |
| 226 | Ok(()) => 0, | |
| 227 | Err(error) => { | |
| 228 | eprintln!("{}: `{command}` failed: {error:#}", user.username); | |
| 229 | let _ = writer.write_all(&git::error_pkt("internal error")).await; | |
| 230 | 1 | |
| 231 | } | |
| 232 | }; | |
| 233 | let _ = writer.shutdown().await; | |
| 234 | let _ = write_half.exit_status(status).await; | |
| 235 | let _ = write_half.eof().await; | |
| 236 | let _ = write_half.close().await; | |
| 237 | }); | |
| 238 | Ok(()) | |
| 239 | } | |
| 240 | } | |
| 241 | ||
| 242 | async fn finish(channel: Channel<Msg>, status: u32) { | |
| 243 | let _ = channel.exit_status(status).await; | |
| 244 | let _ = channel.eof().await; | |
| 245 | let _ = channel.close().await; | |
| 246 | } | |
| 247 | ||
| 248 | /// Parses `git-upload-pack 'owner/repo.git'` and friends. | |
| 249 | fn parse_command(command: &str) -> Option<(Service, &str, &str)> { | |
| 250 | let (program, path) = command.trim().rsplit_once(' ')?; | |
| 251 | let service = match program { | |
| 252 | "git-upload-pack" | "git upload-pack" => Service::UploadPack, | |
| 253 | "git-receive-pack" | "git receive-pack" => Service::ReceivePack, | |
| 254 | _ => return None, | |
| 255 | }; | |
| 256 | let path = path.trim_matches(['\'', '"']).trim_start_matches('/'); | |
| 257 | let path = path.strip_suffix(".git").unwrap_or(path); | |
| 258 | let (owner, repo) = path.split_once('/')?; | |
| 259 | (!owner.is_empty() && !repo.is_empty() && !repo.contains('/')).then_some((service, owner, repo)) | |
| 260 | } | |
| 261 | ||
| 262 | async fn run_git<R, W>( | |
| 263 | api: &Api, | |
| 264 | user: &User, | |
| 265 | command: &str, | |
| 266 | protocol_v2: bool, | |
| 267 | reader: &mut R, | |
| 268 | writer: &mut W, | |
| 269 | ) -> Result<()> | |
| 270 | where | |
| 271 | R: tokio::io::AsyncBufRead + Unpin, | |
| 272 | W: AsyncWrite + Unpin, | |
| 273 | { | |
| 274 | let Some((service, owner, repo)) = parse_command(command) else { | |
| 275 | writer | |
| 276 | .write_all(&git::error_pkt("g1t only supports git over SSH")) | |
| 277 | .await?; | |
| 278 | return Ok(()); | |
| 279 | }; | |
| 280 | match api.access(user, owner, repo, service).await? { | |
| 281 | Ok(access) => git::serve(&api.http, &access, service, protocol_v2, reader, writer).await, | |
| 282 | Err(message) => Ok(writer.write_all(&git::error_pkt(&message)).await?), | |
| 283 | } | |
| 284 | } |