Revise git transport to use I/O more directly
Signed-off-by: Alexis Sellier <alexis@radicle.xyz>
This commit is contained in:
parent
f85a0426d7
commit
48a5c75ae3
|
|
@ -865,48 +865,24 @@ mod tests {
|
||||||
});
|
});
|
||||||
opts.remote_callbacks(callbacks);
|
opts.remote_callbacks(callbacks);
|
||||||
|
|
||||||
let target = git2::Repository::init_bare(target_path).unwrap();
|
|
||||||
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().unwrap();
|
transport::register().unwrap();
|
||||||
|
|
||||||
|
let target = git2::Repository::init_bare(target_path).unwrap();
|
||||||
|
let stream = net::TcpStream::connect(addr).unwrap();
|
||||||
|
let smart = transport::Smart::singleton();
|
||||||
|
|
||||||
|
smart.insert(proj, Box::new(stream.try_clone().unwrap()));
|
||||||
|
|
||||||
// Fetch with the `rad://` transport.
|
// Fetch with the `rad://` transport.
|
||||||
target
|
target
|
||||||
.remote_anonymous(&format!("rad://{}", proj))
|
.remote_anonymous(&format!("rad://{}", proj))
|
||||||
.unwrap()
|
.unwrap()
|
||||||
.fetch(refs, Some(&mut opts), None)
|
.fetch(&["refs/*:refs/*"], Some(&mut opts), None)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
|
||||||
smart.remove(&proj);
|
|
||||||
stream.shutdown(net::Shutdown::Both).unwrap();
|
stream.shutdown(net::Shutdown::Both).unwrap();
|
||||||
|
|
||||||
wt.join().unwrap();
|
|
||||||
rt.join().unwrap();
|
|
||||||
t.join().unwrap();
|
t.join().unwrap();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -13,14 +13,13 @@
|
||||||
//! out the git smart protocol.
|
//! out the git smart protocol.
|
||||||
//!
|
//!
|
||||||
//! This module is meant to be used by first registering our transport with [`register`] and then
|
//! 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`].
|
//! adding or removing streams through [`Smart`], which can be obtained via [`Smart::singleton`].
|
||||||
use std::collections::HashMap;
|
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::sync::{Arc, Mutex};
|
use std::sync::{Arc, Mutex};
|
||||||
|
|
||||||
use crossbeam_channel as chan;
|
use git2::transport::SmartSubtransportStream;
|
||||||
use once_cell::sync::Lazy;
|
use once_cell::sync::Lazy;
|
||||||
|
|
||||||
use crate::git;
|
use crate::git;
|
||||||
|
|
@ -31,6 +30,9 @@ use crate::identity::Id;
|
||||||
/// or its underlying streams.
|
/// or its underlying streams.
|
||||||
static STREAMS: Lazy<Arc<Mutex<HashMap<Id, Stream>>>> = Lazy::new(Default::default);
|
static STREAMS: Lazy<Arc<Mutex<HashMap<Id, Stream>>>> = Lazy::new(Default::default);
|
||||||
|
|
||||||
|
/// The stream associated with a repository.
|
||||||
|
type Stream = Box<dyn SmartSubtransportStream>;
|
||||||
|
|
||||||
/// Git transport protocol over an I/O stream.
|
/// Git transport protocol over an I/O stream.
|
||||||
#[derive(Clone)]
|
#[derive(Clone)]
|
||||||
pub struct Smart {
|
pub struct Smart {
|
||||||
|
|
@ -39,17 +41,23 @@ pub struct Smart {
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Smart {
|
impl Smart {
|
||||||
pub fn get(&self, id: &Id) -> Option<Stream> {
|
/// Get access to the radicle smart transport protocol.
|
||||||
self.streams.lock().unwrap().get(id).cloned()
|
/// The returned object has mutable access to the underlying stream map, and is safe to clone.
|
||||||
|
pub fn singleton() -> Self {
|
||||||
|
Self {
|
||||||
|
streams: STREAMS.clone(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Take a stream from the map.
|
||||||
|
/// This makes the stream unavailable until it is re-inserted.
|
||||||
|
pub fn take(&self, id: &Id) -> Option<Stream> {
|
||||||
|
self.streams.lock().unwrap().remove(id)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn insert(&self, id: Id, stream: Stream) {
|
pub fn insert(&self, id: Id, stream: Stream) {
|
||||||
self.streams.lock().unwrap().insert(id, 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 {
|
||||||
|
|
@ -74,7 +82,7 @@ impl git2::transport::SmartSubtransport for Smart {
|
||||||
return Err(git2::Error::from_str("Git URL scheme must be `rad`"));
|
return Err(git2::Error::from_str("Git URL scheme must be `rad`"));
|
||||||
}
|
}
|
||||||
|
|
||||||
if let Some(stream) = self.get(&id) {
|
if let Some(stream) = self.take(&id) {
|
||||||
match action {
|
match action {
|
||||||
git2::transport::Service::UploadPackLs => {}
|
git2::transport::Service::UploadPackLs => {}
|
||||||
git2::transport::Service::UploadPack => {}
|
git2::transport::Service::UploadPack => {}
|
||||||
|
|
@ -89,7 +97,7 @@ impl git2::transport::SmartSubtransport for Smart {
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Ok(Box::new(stream))
|
Ok(stream)
|
||||||
} else {
|
} else {
|
||||||
Err(git2::Error::from_str(&format!(
|
Err(git2::Error::from_str(&format!(
|
||||||
"repository {id} does not have an associated stream"
|
"repository {id} does not have an associated stream"
|
||||||
|
|
@ -102,60 +110,6 @@ impl git2::transport::SmartSubtransport for Smart {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// A byte stream connected to some I/O source.
|
|
||||||
/// 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 {
|
|
||||||
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
|
|
||||||
self.send
|
|
||||||
.send(buf.to_owned())
|
|
||||||
.map_err(|e| io::Error::new(io::ErrorKind::BrokenPipe, e))?;
|
|
||||||
|
|
||||||
Ok(buf.len())
|
|
||||||
}
|
|
||||||
|
|
||||||
fn flush(&mut self) -> io::Result<()> {
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
impl io::Read for Stream {
|
|
||||||
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
|
|
||||||
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 radicle transport with `git`.
|
/// Register the radicle transport with `git`.
|
||||||
///
|
///
|
||||||
/// Returns an error if called more than once.
|
/// Returns an error if called more than once.
|
||||||
|
|
@ -168,7 +122,7 @@ pub fn register() -> Result<(), git2::Error> {
|
||||||
unsafe {
|
unsafe {
|
||||||
let prefix = git::url::Scheme::Radicle.to_string();
|
let prefix = git::url::Scheme::Radicle.to_string();
|
||||||
git2::transport::register(&prefix, move |remote| {
|
git2::transport::register(&prefix, move |remote| {
|
||||||
git2::transport::Transport::smart(remote, false, self::smart())
|
git2::transport::Transport::smart(remote, false, Smart::singleton())
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
|
|
@ -177,13 +131,3 @@ pub fn register() -> 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