Flesh out custom git transport
Signed-off-by: Alexis Sellier <alexis@radicle.xyz>
This commit is contained in:
parent
b4dcc00fe3
commit
1da38708fc
|
|
@ -768,6 +768,7 @@ dependencies = [
|
||||||
name = "radicle"
|
name = "radicle"
|
||||||
version = "0.2.0"
|
version = "0.2.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
|
"crossbeam-channel",
|
||||||
"ed25519-compact",
|
"ed25519-compact",
|
||||||
"fastrand",
|
"fastrand",
|
||||||
"git-ref-format",
|
"git-ref-format",
|
||||||
|
|
|
||||||
|
|
@ -10,6 +10,7 @@ default = []
|
||||||
test = ["quickcheck"]
|
test = ["quickcheck"]
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
|
crossbeam-channel = { version = "0.5.6" }
|
||||||
ed25519-compact = { version = "1.0.12", features = ["pem"] }
|
ed25519-compact = { version = "1.0.12", features = ["pem"] }
|
||||||
fastrand = { version = "1.8.0" }
|
fastrand = { version = "1.8.0" }
|
||||||
git-ref-format = { version = "0", features = ["serde", "macro"] }
|
git-ref-format = { version = "0", features = ["serde", "macro"] }
|
||||||
|
|
|
||||||
|
|
@ -652,9 +652,13 @@ pub mod paths {
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
|
use std::io::{Read, Write};
|
||||||
|
use std::{io, net, process, thread};
|
||||||
|
|
||||||
use super::*;
|
use super::*;
|
||||||
use crate::assert_matches;
|
use crate::assert_matches;
|
||||||
use crate::git;
|
use crate::git;
|
||||||
|
use crate::rad;
|
||||||
use crate::storage::refs::SIGNATURE_REF;
|
use crate::storage::refs::SIGNATURE_REF;
|
||||||
use crate::storage::{ReadRepository, ReadStorage, RefUpdate, WriteRepository};
|
use crate::storage::{ReadRepository, ReadStorage, RefUpdate, WriteRepository};
|
||||||
use crate::test::arbitrary;
|
use crate::test::arbitrary;
|
||||||
|
|
@ -793,23 +797,33 @@ mod tests {
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_upload_pack() {
|
fn test_upload_pack() {
|
||||||
use std::io::{Read, Write};
|
let tmp = tempfile::tempdir().unwrap();
|
||||||
use std::{io, net, process, thread};
|
let signer = MockSigner::default();
|
||||||
|
let remote = *signer.public_key();
|
||||||
|
let storage = Storage::open(tmp.path().join("storage")).unwrap();
|
||||||
let socket = net::TcpListener::bind(net::SocketAddr::from(([0, 0, 0, 0], 0))).unwrap();
|
let socket = net::TcpListener::bind(net::SocketAddr::from(([0, 0, 0, 0], 0))).unwrap();
|
||||||
let addr = socket.local_addr().unwrap();
|
let addr = socket.local_addr().unwrap();
|
||||||
let tmp = tempfile::tempdir().unwrap();
|
|
||||||
let source_path = tmp.path().join("source");
|
let source_path = tmp.path().join("source");
|
||||||
let target_path = tmp.path().join("target");
|
let target_path = tmp.path().join("target");
|
||||||
let (_source, _) = fixtures::repository(&source_path);
|
let (source, _) = fixtures::repository(&source_path);
|
||||||
|
let (proj, _) = rad::init(
|
||||||
|
&source,
|
||||||
|
"radicle",
|
||||||
|
"radicle",
|
||||||
|
BranchName::from("master"),
|
||||||
|
signer,
|
||||||
|
&storage,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
let t = thread::spawn(move || {
|
let t = thread::spawn(move || {
|
||||||
let (stream, _) = socket.accept().unwrap();
|
let (stream, _) = socket.accept().unwrap();
|
||||||
// NOTE: `--stateless-rpc` doesn't work, neither does `GIT_PROTOCOL=version=2`.
|
let repo = storage.repository(proj).unwrap();
|
||||||
|
// NOTE: `GIT_PROTOCOL=version=2` doesn't work.
|
||||||
let mut child = process::Command::new("git")
|
let mut child = process::Command::new("git")
|
||||||
.current_dir(source_path.join(".git"))
|
.current_dir(repo.path())
|
||||||
.arg("upload-pack")
|
.arg("upload-pack")
|
||||||
.arg("--strict")
|
.arg("--strict") // The path to the git repo must be exact.
|
||||||
.arg(".")
|
.arg(".")
|
||||||
.stdout(process::Stdio::piped())
|
.stdout(process::Stdio::piped())
|
||||||
.stdin(process::Stdio::piped())
|
.stdin(process::Stdio::piped())
|
||||||
|
|
@ -854,18 +868,56 @@ mod tests {
|
||||||
let target = git2::Repository::init_bare(target_path).unwrap();
|
let target = git2::Repository::init_bare(target_path).unwrap();
|
||||||
let refs: &[&str] = &["refs/*:refs/*"];
|
let refs: &[&str] = &["refs/*:refs/*"];
|
||||||
|
|
||||||
|
let stream = net::TcpStream::connect(addr).unwrap();
|
||||||
|
let mut stream_r = stream.try_clone().unwrap();
|
||||||
|
let mut stream_w = stream.try_clone().unwrap();
|
||||||
|
let (stream_r_send, stream_r_recv) = crossbeam_channel::unbounded::<Vec<u8>>();
|
||||||
|
let (stream_w_send, stream_w_recv) = crossbeam_channel::unbounded::<Vec<u8>>();
|
||||||
|
|
||||||
|
let smart = transport::smart();
|
||||||
|
smart.insert(proj, transport::Stream::new(stream_w_send, stream_r_recv));
|
||||||
|
|
||||||
|
let rt = thread::spawn(move || {
|
||||||
|
let mut buf = vec![0u8; 1024];
|
||||||
|
|
||||||
|
while let Ok(n) = stream_r.read(&mut buf) {
|
||||||
|
if n == 0 {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
stream_r_send.send(buf[..n].to_vec()).unwrap();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
let wt = thread::spawn(move || {
|
||||||
|
for buf in stream_w_recv.iter() {
|
||||||
|
stream_w.write_all(&buf).unwrap();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
// Register the `rad://` transport.
|
// Register the `rad://` transport.
|
||||||
transport::register("rad").unwrap();
|
transport::register("rad").unwrap();
|
||||||
// Fetch with the `rad://` transport.
|
// Fetch with the `rad://` transport.
|
||||||
target
|
target
|
||||||
.remote_anonymous(&format!("rad://{}", addr))
|
.remote_anonymous(&format!("rad://{}", proj))
|
||||||
.unwrap()
|
.unwrap()
|
||||||
.fetch(refs, Some(&mut opts), None)
|
.fetch(refs, Some(&mut opts), None)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
|
||||||
|
smart.remove(&proj);
|
||||||
|
stream.shutdown(net::Shutdown::Both).unwrap();
|
||||||
|
|
||||||
|
wt.join().unwrap();
|
||||||
|
rt.join().unwrap();
|
||||||
t.join().unwrap();
|
t.join().unwrap();
|
||||||
}
|
}
|
||||||
assert_eq!(updates, vec![String::from("refs/heads/master")]);
|
|
||||||
|
assert_eq!(
|
||||||
|
updates,
|
||||||
|
vec![
|
||||||
|
format!("refs/remotes/{remote}/heads/master"),
|
||||||
|
format!("refs/remotes/{remote}/heads/radicle/id"),
|
||||||
|
format!("refs/remotes/{remote}/radicle/signature")
|
||||||
|
]
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
|
|
|
||||||
|
|
@ -1,13 +1,65 @@
|
||||||
|
//! Git sub-transport used for fetching radicle data.
|
||||||
|
//!
|
||||||
|
//! To have control over the communication, and to allow git streams to be multiplexed over
|
||||||
|
//! existing TCP connections, we implement the [`git2::transport::SmartSubtransport`] trait.
|
||||||
|
//!
|
||||||
|
//! We choose `rad` as the URL scheme for this custom transport, and include only the identity
|
||||||
|
//! of the repository we're looking to fetch, eg. `rad://zP1GztjSdYNHK7jpdrXbaJ6Ki2Ke`, since
|
||||||
|
//! we expect a connection to a host to already be established.
|
||||||
|
//!
|
||||||
|
//! We then maintain a map from identifier to stream, for all active streams, ie. streams that
|
||||||
|
//! are associated with an underlying TCP connection. When a URL is requested, we lookup
|
||||||
|
//! the stream and return it to the [`git2`] smart-protocol implementation, so that it can carry
|
||||||
|
//! out the git smart protocol.
|
||||||
|
//!
|
||||||
|
//! This module is meant to be used by first registering our transport with [`register`] and then
|
||||||
|
//! adding or removing streams through [`Smart`], which can be obtained by calling [`smart`].
|
||||||
|
use std::collections::HashMap;
|
||||||
|
use std::io;
|
||||||
use std::str::FromStr;
|
use std::str::FromStr;
|
||||||
use std::sync::atomic;
|
use std::sync::atomic;
|
||||||
use std::{io, net};
|
use std::sync::{Arc, Mutex};
|
||||||
|
|
||||||
|
use crossbeam_channel as chan;
|
||||||
|
use once_cell::sync::Lazy;
|
||||||
|
|
||||||
use crate::git;
|
use crate::git;
|
||||||
|
use crate::identity::Id;
|
||||||
|
|
||||||
/// Git smart protocol over a TCP stream.
|
/// The map of git smart sub-transport streams. We keep a global map because we have
|
||||||
pub struct Smart;
|
/// no control over how [`git2::transport::register`] instantiates our [`Smart`] transport
|
||||||
|
/// or its underlying streams.
|
||||||
|
static STREAMS: Lazy<Arc<Mutex<HashMap<Id, Stream>>>> = Lazy::new(Default::default);
|
||||||
|
|
||||||
|
/// Git transport protocol over an I/O stream.
|
||||||
|
#[derive(Clone)]
|
||||||
|
pub struct Smart {
|
||||||
|
/// The underlying active streams, keyed by repository identifier.
|
||||||
|
streams: Arc<Mutex<HashMap<Id, Stream>>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Smart {
|
||||||
|
pub fn get(&self, id: &Id) -> Option<Stream> {
|
||||||
|
self.streams.lock().unwrap().get(id).cloned()
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn insert(&self, id: Id, stream: Stream) {
|
||||||
|
self.streams.lock().unwrap().insert(id, stream);
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn remove(&self, id: &Id) {
|
||||||
|
self.streams.lock().unwrap().remove(id);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
impl git2::transport::SmartSubtransport for Smart {
|
impl git2::transport::SmartSubtransport for Smart {
|
||||||
|
/// Run a git service on this transport.
|
||||||
|
///
|
||||||
|
/// Based on the URL, which must be of the form `rad://zP1GztjSdYNHK7jpdrXbaJ6Ki2Ke`,
|
||||||
|
/// we retrieve an underlying stream and return it.
|
||||||
|
///
|
||||||
|
/// We only support the upload-pack service, since only fetches are authorized by the
|
||||||
|
/// remote.
|
||||||
fn action(
|
fn action(
|
||||||
&self,
|
&self,
|
||||||
url: &str,
|
url: &str,
|
||||||
|
|
@ -15,36 +67,34 @@ impl git2::transport::SmartSubtransport for Smart {
|
||||||
) -> Result<Box<dyn git2::transport::SmartSubtransportStream>, git2::Error> {
|
) -> Result<Box<dyn git2::transport::SmartSubtransportStream>, git2::Error> {
|
||||||
let url = git::Url::from_bytes(url.as_bytes())
|
let url = git::Url::from_bytes(url.as_bytes())
|
||||||
.map_err(|e| git2::Error::from_str(e.to_string().as_str()))?;
|
.map_err(|e| git2::Error::from_str(e.to_string().as_str()))?;
|
||||||
|
let id = Id::from_str(url.host.unwrap_or_default().as_str())
|
||||||
|
.map_err(|_| git2::Error::from_str("Git URL does not contain a valid project id"))?;
|
||||||
|
|
||||||
let addr = if let (Some(host), Some(port)) = (url.host, url.port) {
|
if url.scheme != git::url::Scheme::Radicle {
|
||||||
// TODO: Support hostnames.
|
return Err(git2::Error::from_str("Git URL scheme must be `rad`"));
|
||||||
net::SocketAddr::new(
|
}
|
||||||
net::IpAddr::from_str(&host)
|
|
||||||
.map_err(|e| git2::Error::from_str(e.to_string().as_str()))?,
|
if let Some(stream) = self.get(&id) {
|
||||||
port,
|
match action {
|
||||||
)
|
git2::transport::Service::UploadPackLs => {}
|
||||||
} else {
|
git2::transport::Service::UploadPack => {}
|
||||||
return Err(git2::Error::from_str("Git URL must have a host and port"));
|
git2::transport::Service::ReceivePack => {
|
||||||
};
|
return Err(git2::Error::from_str(
|
||||||
|
"git-receive-pack is not supported with the custom transport",
|
||||||
let stream = std::net::TcpStream::connect(addr)
|
));
|
||||||
.map_err(|e| git2::Error::from_str(e.to_string().as_str()))?;
|
}
|
||||||
|
git2::transport::Service::ReceivePackLs => {
|
||||||
match action {
|
return Err(git2::Error::from_str(
|
||||||
git2::transport::Service::UploadPackLs => {}
|
"git-receive-pack is not supported with the custom transport",
|
||||||
git2::transport::Service::UploadPack => {}
|
));
|
||||||
git2::transport::Service::ReceivePack => {
|
}
|
||||||
return Err(git2::Error::from_str(
|
}
|
||||||
"git-receive-pack is not supported with the custom transport",
|
Ok(Box::new(stream))
|
||||||
));
|
} else {
|
||||||
}
|
Err(git2::Error::from_str(&format!(
|
||||||
git2::transport::Service::ReceivePackLs => {
|
"repository {id} does not have an associated stream"
|
||||||
return Err(git2::Error::from_str(
|
)))
|
||||||
"git-receive-pack is not supported with the custom transport",
|
|
||||||
));
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
Ok(Box::new(Stream { stream }))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn close(&self) -> Result<(), git2::Error> {
|
fn close(&self) -> Result<(), git2::Error> {
|
||||||
|
|
@ -52,34 +102,72 @@ impl git2::transport::SmartSubtransport for Smart {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
struct Stream {
|
/// A byte stream connected to some I/O source.
|
||||||
stream: std::net::TcpStream,
|
/// One of these is created for every git operation, eg. `fetch`.
|
||||||
|
#[derive(Clone)]
|
||||||
|
pub struct Stream {
|
||||||
|
/// Send bytes to the network.
|
||||||
|
send: chan::Sender<Vec<u8>>,
|
||||||
|
/// Receive bytes from the network.
|
||||||
|
recv: chan::Receiver<Vec<u8>>,
|
||||||
|
/// Bytes read from the receive channel that didn't fit in the read buffer.
|
||||||
|
pending: Vec<u8>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Stream {
|
||||||
|
/// Create a new stream from a sender and receiver.
|
||||||
|
pub fn new(send: chan::Sender<Vec<u8>>, recv: chan::Receiver<Vec<u8>>) -> Self {
|
||||||
|
Self {
|
||||||
|
send,
|
||||||
|
recv,
|
||||||
|
pending: Vec::new(),
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl io::Write for Stream {
|
impl io::Write for Stream {
|
||||||
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
|
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
|
||||||
self.stream.write(buf)
|
self.send
|
||||||
|
.send(buf.to_owned())
|
||||||
|
.map_err(|e| io::Error::new(io::ErrorKind::BrokenPipe, e))?;
|
||||||
|
|
||||||
|
Ok(buf.len())
|
||||||
}
|
}
|
||||||
|
|
||||||
fn flush(&mut self) -> std::io::Result<()> {
|
fn flush(&mut self) -> io::Result<()> {
|
||||||
self.stream.flush()
|
Ok(())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl io::Read for Stream {
|
impl io::Read for Stream {
|
||||||
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
|
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
|
||||||
self.stream.read(buf)
|
let bytes = self
|
||||||
|
.recv
|
||||||
|
.recv()
|
||||||
|
.map_err(|e| io::Error::new(io::ErrorKind::BrokenPipe, e))?;
|
||||||
|
self.pending.extend(&bytes);
|
||||||
|
|
||||||
|
// There must be a nicer way to do this...
|
||||||
|
let count = buf.len().min(self.pending.len());
|
||||||
|
buf[..count].copy_from_slice(&self.pending[..count]);
|
||||||
|
self.pending.drain(..count);
|
||||||
|
|
||||||
|
Ok(count)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Register the "smart" transport with `git`.
|
/// Register the radicle transport with `git`.
|
||||||
|
///
|
||||||
|
/// Returns an error if called more than once.
|
||||||
|
///
|
||||||
pub fn register(prefix: &str) -> Result<(), git2::Error> {
|
pub fn register(prefix: &str) -> Result<(), git2::Error> {
|
||||||
static REGISTERED: atomic::AtomicBool = atomic::AtomicBool::new(false);
|
static REGISTERED: atomic::AtomicBool = atomic::AtomicBool::new(false);
|
||||||
|
|
||||||
|
// Registration is not thread-safe, so make sure we prevent re-entrancy.
|
||||||
if !REGISTERED.swap(true, atomic::Ordering::SeqCst) {
|
if !REGISTERED.swap(true, atomic::Ordering::SeqCst) {
|
||||||
unsafe {
|
unsafe {
|
||||||
git2::transport::register(prefix, move |remote| {
|
git2::transport::register(prefix, move |remote| {
|
||||||
git2::transport::Transport::smart(remote, false, Smart)
|
git2::transport::Transport::smart(remote, false, self::smart())
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
|
|
@ -88,3 +176,13 @@ pub fn register(prefix: &str) -> Result<(), git2::Error> {
|
||||||
))
|
))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Get access to the radicle smart transport protocol.
|
||||||
|
///
|
||||||
|
/// The returned object has mutable access to the underlying stream map, and is safe to clone.
|
||||||
|
///
|
||||||
|
pub fn smart() -> Smart {
|
||||||
|
Smart {
|
||||||
|
streams: STREAMS.clone(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue