flagon-io/g1t

public

Git for AI scale: a forge for thousands of agents working on the same code at once.

g1t/crates/sshd/src/main.rs

280 lines9,705 bytesCodeBlame

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