Consolidate `node::Handle` traits into one

Signed-off-by: Alexis Sellier <alexis@radicle.xyz>
This commit is contained in:
Alexis Sellier 2022-12-07 15:34:46 +01:00
parent a100f1c683
commit 7b8bb08b08
No known key found for this signature in database
11 changed files with 130 additions and 99 deletions

View File

@ -80,12 +80,12 @@ pub fn run(options: Options, ctx: impl term::Context) -> anyhow::Result<()> {
pub fn clone(id: Id, _interactive: Interactive, ctx: impl term::Context) -> anyhow::Result<()> { pub fn clone(id: Id, _interactive: Interactive, ctx: impl term::Context) -> anyhow::Result<()> {
let profile = ctx.profile()?; let profile = ctx.profile()?;
let node = radicle::node::connect(profile.node())?; let mut node = radicle::node::connect(profile.node())?;
let signer = term::signer(&profile)?; let signer = term::signer(&profile)?;
// Track & fetch project. // Track & fetch project.
node.track_repo(&id).context("track")?; node.track_repo(id).context("track")?;
node.fetch(&id).context("fetch")?; node.fetch(id).context("fetch")?;
// Create a local fork of the project, under our own id. // Create a local fork of the project, under our own id.
rad::fork(id, &signer, &profile.storage).context("fork")?; rad::fork(id, &signer, &profile.storage).context("fork")?;

View File

@ -95,7 +95,7 @@ pub fn run(options: Options, ctx: impl term::Context) -> anyhow::Result<()> {
let storage = &profile.storage; let storage = &profile.storage;
let (_, rid) = radicle::rad::cwd().context("this command must be run within a project")?; let (_, rid) = radicle::rad::cwd().context("this command must be run within a project")?;
let Doc { payload, .. } = storage.repository(rid)?.project_of(profile.id())?; let Doc { payload, .. } = storage.repository(rid)?.project_of(profile.id())?;
let node = radicle::node::connect(&profile.node())?; let mut node = radicle::node::connect(&profile.node())?;
term::info!( term::info!(
"Establishing 🌱 tracking relationship for {}", "Establishing 🌱 tracking relationship for {}",
@ -103,7 +103,7 @@ pub fn run(options: Options, ctx: impl term::Context) -> anyhow::Result<()> {
); );
term::blank(); term::blank();
let tracked = node.track_node(&peer, options.alias.as_deref())?; let tracked = node.track_node(peer, options.alias.clone())?;
let outcome = if tracked { "established" } else { "exists" }; let outcome = if tracked { "established" } else { "exists" };
if let Some(alias) = options.alias { if let Some(alias) = options.alias {
@ -118,7 +118,7 @@ pub fn run(options: Options, ctx: impl term::Context) -> anyhow::Result<()> {
} }
if options.fetch { if options.fetch {
node.fetch(&rid)?; node.fetch(rid)?;
} }
Ok(()) Ok(())

View File

@ -88,6 +88,6 @@ pub fn run(options: Options, ctx: impl term::Context) -> anyhow::Result<()> {
} }
pub fn untrack(id: Id, profile: &Profile) -> anyhow::Result<bool> { pub fn untrack(id: Id, profile: &Profile) -> anyhow::Result<bool> {
let node = radicle::node::connect(profile.node())?; let mut node = radicle::node::connect(profile.node())?;
node.untrack_repo(&id).map_err(|e| anyhow!(e)) node.untrack_repo(id).map_err(|e| anyhow!(e))
} }

View File

@ -55,12 +55,25 @@ pub struct Handle<W: Waker> {
pub(crate) waker: W, pub(crate) waker: W,
} }
impl<W: Waker> traits::Handle for Handle<W> { impl<W: Waker> Handle<W> {
fn command(&self, cmd: service::Command) -> Result<(), Error> {
self.commands.send(cmd)?;
self.waker.wake()?;
Ok(())
}
}
impl<W: Waker> radicle::node::Handle for Handle<W> {
type Session = Session;
type FetchLookup = FetchLookup;
type Error = Error;
fn listening(&self) -> Result<net::SocketAddr, Error> { fn listening(&self) -> Result<net::SocketAddr, Error> {
self.listening.recv().map_err(Error::from) self.listening.recv().map_err(Error::from)
} }
fn fetch(&mut self, id: Id) -> Result<FetchLookup, Error> { fn fetch(&mut self, id: Id) -> Result<Self::FetchLookup, Error> {
let (sender, receiver) = chan::bounded(1); let (sender, receiver) = chan::bounded(1);
self.commands.send(service::Command::Fetch(id, sender))?; self.commands.send(service::Command::Fetch(id, sender))?;
receiver.recv().map_err(Error::from) receiver.recv().map_err(Error::from)
@ -98,13 +111,6 @@ impl<W: Waker> traits::Handle for Handle<W> {
self.command(service::Command::AnnounceRefs(id)) self.command(service::Command::AnnounceRefs(id))
} }
fn command(&self, cmd: service::Command) -> Result<(), Error> {
self.commands.send(cmd)?;
self.waker.wake()?;
Ok(())
}
fn routing(&self) -> Result<chan::Receiver<(Id, NodeId)>, Error> { fn routing(&self) -> Result<chan::Receiver<(Id, NodeId)>, Error> {
let (sender, receiver) = chan::unbounded(); let (sender, receiver) = chan::unbounded();
let query: Arc<QueryState> = Arc::new(move |state| { let query: Arc<QueryState> = Arc::new(move |state| {
@ -151,35 +157,3 @@ impl<W: Waker> traits::Handle for Handle<W> {
Ok(()) Ok(())
} }
} }
pub mod traits {
use super::*;
pub trait Handle {
/// Wait for the node's listening socket to be bound.
fn listening(&self) -> Result<net::SocketAddr, Error>;
/// Retrieve or update the project from network.
fn fetch(&mut self, id: Id) -> Result<FetchLookup, Error>;
/// Start tracking the given project. Doesn't do anything if the project is already
/// tracked.
fn track_repo(&mut self, id: Id) -> Result<bool, Error>;
/// Start tracking the given node.
fn track_node(&mut self, id: NodeId, alias: Option<String>) -> Result<bool, Error>;
/// Untrack the given project and delete it from storage.
fn untrack_repo(&mut self, id: Id) -> Result<bool, Error>;
/// Untrack the given node.
fn untrack_node(&mut self, id: NodeId) -> Result<bool, Error>;
/// Notify the client that a project has been updated.
fn announce_refs(&mut self, id: Id) -> Result<(), Error>;
/// Send a command to the command channel, and wake up the event loop.
fn command(&self, cmd: service::Command) -> Result<(), Error>;
/// Ask the client to shutdown.
fn shutdown(self) -> Result<(), Error>;
/// Query the routing table entries.
fn routing(&self) -> Result<chan::Receiver<(Id, NodeId)>, Error>;
/// Query the peer session state.
fn sessions(&self) -> Result<chan::Receiver<(NodeId, Session)>, Error>;
/// Query the inventory.
fn inventory(&self) -> Result<chan::Receiver<Id>, Error>;
}
}

View File

@ -7,8 +7,9 @@ use std::os::unix::net::UnixStream;
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
use std::{fs, io, net}; use std::{fs, io, net};
use radicle::node::Handle;
use crate::client; use crate::client;
use crate::client::handle::traits::Handle;
use crate::identity::Id; use crate::identity::Id;
use crate::node; use crate::node;
use crate::service::FetchLookup; use crate::service::FetchLookup;
@ -23,7 +24,13 @@ pub enum Error {
} }
/// Listen for commands on the control socket, and process them. /// Listen for commands on the control socket, and process them.
pub fn listen<P: AsRef<Path>, H: Handle>(path: P, mut handle: H) -> Result<(), Error> { pub fn listen<
P: AsRef<Path>,
H: Handle<Error = client::handle::Error, FetchLookup = FetchLookup>,
>(
path: P,
mut handle: H,
) -> Result<(), Error> {
// Remove the socket file on startup before rebinding. // Remove the socket file on startup before rebinding.
fs::remove_file(&path).ok(); fs::remove_file(&path).ok();
fs::create_dir_all( fs::create_dir_all(
@ -69,7 +76,10 @@ enum DrainError {
Io(#[from] io::Error), Io(#[from] io::Error),
} }
fn drain<H: Handle>(stream: &UnixStream, handle: &mut H) -> Result<(), DrainError> { fn drain<H: Handle<Error = client::handle::Error, FetchLookup = FetchLookup>>(
stream: &UnixStream,
handle: &mut H,
) -> Result<(), DrainError> {
let mut reader = BufReader::new(stream); let mut reader = BufReader::new(stream);
let mut writer = LineWriter::new(stream); let mut writer = LineWriter::new(stream);
@ -198,7 +208,11 @@ fn drain<H: Handle>(stream: &UnixStream, handle: &mut H) -> Result<(), DrainErro
Ok(()) Ok(())
} }
fn fetch<W: Write, H: Handle>(id: Id, mut writer: W, handle: &mut H) -> Result<(), DrainError> { fn fetch<W: Write, H: Handle<Error = client::handle::Error, FetchLookup = FetchLookup>>(
id: Id,
mut writer: W,
handle: &mut H,
) -> Result<(), DrainError> {
match handle.fetch(id) { match handle.fetch(id) {
Err(e) => { Err(e) => {
return Err(DrainError::Client(e)); return Err(DrainError::Client(e));
@ -305,20 +319,24 @@ mod tests {
move || crate::control::listen(socket, handle) move || crate::control::listen(socket, handle)
}); });
let handle = loop { let mut handle = loop {
if let Ok(conn) = Node::connect(&socket) { if let Ok(conn) = Node::connect(&socket) {
break conn; break conn;
} }
}; };
assert!(handle.track_repo(&proj).unwrap()); assert!(handle.track_repo(proj).unwrap());
assert!(!handle.track_repo(&proj).unwrap()); assert!(!handle.track_repo(proj).unwrap());
assert!(handle.untrack_repo(&proj).unwrap()); assert!(handle.untrack_repo(proj).unwrap());
assert!(!handle.untrack_repo(&proj).unwrap()); assert!(!handle.untrack_repo(proj).unwrap());
assert!(handle.track_node(&peer, Some("alice")).unwrap()); assert!(handle
assert!(!handle.track_node(&peer, Some("alice")).unwrap()); .track_node(peer, Some(String::from("alice")))
assert!(handle.untrack_node(&peer).unwrap()); .unwrap());
assert!(!handle.untrack_node(&peer).unwrap()); assert!(!handle
.track_node(peer, Some(String::from("alice")))
.unwrap());
assert!(handle.untrack_node(peer).unwrap());
assert!(!handle.untrack_node(peer).unwrap());
} }
} }

View File

@ -3,7 +3,6 @@ use std::sync::{Arc, Mutex};
use crossbeam_channel as chan; use crossbeam_channel as chan;
use crate::client::handle::traits;
use crate::client::handle::Error; use crate::client::handle::Error;
use crate::identity::Id; use crate::identity::Id;
use crate::service; use crate::service;
@ -17,7 +16,11 @@ pub struct Handle {
pub tracking_nodes: HashSet<NodeId>, pub tracking_nodes: HashSet<NodeId>,
} }
impl traits::Handle for Handle { impl radicle::node::Handle for Handle {
type Error = Error;
type Session = service::Session;
type FetchLookup = FetchLookup;
fn listening(&self) -> Result<std::net::SocketAddr, Error> { fn listening(&self) -> Result<std::net::SocketAddr, Error> {
unimplemented!() unimplemented!()
} }
@ -48,10 +51,6 @@ impl traits::Handle for Handle {
Ok(()) Ok(())
} }
fn command(&self, _cmd: service::Command) -> Result<(), Error> {
Ok(())
}
fn routing(&self) -> Result<chan::Receiver<(Id, service::NodeId)>, Error> { fn routing(&self) -> Result<chan::Receiver<(Id, service::NodeId)>, Error> {
unimplemented!(); unimplemented!();
} }

View File

@ -108,8 +108,8 @@ pub fn run(profile: radicle::Profile) -> Result<(), Box<dyn std::error::Error +
// Connect to local node and announce refs to the network. // Connect to local node and announce refs to the network.
// If our node is not running, we simply skip this step, as the // If our node is not running, we simply skip this step, as the
// refs will be announced eventually, when the node restarts. // refs will be announced eventually, when the node restarts.
if let Ok(conn) = radicle::node::connect(&profile.node()) { if let Ok(mut conn) = radicle::node::connect(&profile.node()) {
conn.announce_refs(&url.repo)?; conn.announce_refs(url.repo)?;
} }
} }
} }

View File

@ -11,8 +11,8 @@ fn main() -> anyhow::Result<()> {
if let Some(id) = env::args().nth(1) { if let Some(id) = env::args().nth(1) {
let id = Id::from_str(&id)?; let id = Id::from_str(&id)?;
let node = radicle::node::connect(profile.node())?; let mut node = radicle::node::connect(profile.node())?;
let repo = radicle::rad::clone(id, &cwd, &signer, &profile.storage, &node)?; let repo = radicle::rad::clone(id, &cwd, &signer, &profile.storage, &mut node)?;
println!( println!(
"ok: project {id} cloned into `{}`", "ok: project {id} cloned into `{}`",

View File

@ -16,7 +16,7 @@ fn main() -> anyhow::Result<()> {
let sigrefs = project.sign_refs(&signer)?; let sigrefs = project.sign_refs(&signer)?;
let head = project.set_head()?; let head = project.set_head()?;
radicle::node::connect(&profile.node())?.announce_refs(&id)?; radicle::node::connect(&profile.node())?.announce_refs(id)?;
println!("head: {}", head); println!("head: {}", head);
println!("ok: {}", sigrefs.signature); println!("ok: {}", sigrefs.signature);

View File

@ -1,12 +1,13 @@
mod features; mod features;
use std::io;
use std::io::{BufRead, BufReader, Write}; use std::io::{BufRead, BufReader, Write};
use std::os::unix::net::UnixStream; use std::os::unix::net::UnixStream;
use std::path::Path; use std::path::Path;
use std::{io, net};
use crate::crypto::PublicKey; use crate::crypto::PublicKey;
use crate::identity::Id; use crate::identity::Id;
use crossbeam_channel as chan;
pub use features::Features; pub use features::Features;
@ -27,22 +28,38 @@ pub enum Error {
EmptyResponse { cmd: &'static str }, EmptyResponse { cmd: &'static str },
} }
/// A handle to send commands to the node or request information.
pub trait Handle { pub trait Handle {
/// Fetch a project from the network. Fails if the project isn't tracked. /// The result of a fetch request.
fn fetch(&self, id: &Id) -> Result<(), Error>; type FetchLookup;
/// Start tracking the given node. If the node is already tracked, /// The peer session type.
/// updates the alias if necessary. type Session;
fn track_node(&self, id: &NodeId, alias: Option<&str>) -> Result<bool, Error>; /// The error returned by all methods.
/// Start tracking the given repository. type Error: std::error::Error;
fn track_repo(&self, id: &Id) -> Result<bool, Error>;
/// Wait for the node's listening socket to be bound.
fn listening(&self) -> Result<net::SocketAddr, Self::Error>;
/// Retrieve or update the project from network.
fn fetch(&mut self, id: Id) -> Result<Self::FetchLookup, Self::Error>;
/// Start tracking the given project. Doesn't do anything if the project is already
/// tracked.
fn track_repo(&mut self, id: Id) -> Result<bool, Self::Error>;
/// Start tracking the given node.
fn track_node(&mut self, id: NodeId, alias: Option<String>) -> Result<bool, Self::Error>;
/// Untrack the given project and delete it from storage.
fn untrack_repo(&mut self, id: Id) -> Result<bool, Self::Error>;
/// Untrack the given node. /// Untrack the given node.
fn untrack_node(&self, id: &NodeId) -> Result<bool, Error>; fn untrack_node(&mut self, id: NodeId) -> Result<bool, Self::Error>;
/// Untrack the given repository and delete it from storage. /// Notify the client that a project has been updated.
fn untrack_repo(&self, id: &Id) -> Result<bool, Error>; fn announce_refs(&mut self, id: Id) -> Result<(), Self::Error>;
/// Notify the network that we have new refs. /// Ask the client to shutdown.
fn announce_refs(&self, id: &Id) -> Result<(), Error>; fn shutdown(self) -> Result<(), Self::Error>;
/// Ask the node to shutdown. /// Query the routing table entries.
fn shutdown(self) -> Result<(), Error>; fn routing(&self) -> Result<chan::Receiver<(Id, NodeId)>, Self::Error>;
/// Query the peer session state.
fn sessions(&self) -> Result<chan::Receiver<(NodeId, Self::Session)>, Self::Error>;
/// Query the inventory.
fn inventory(&self) -> Result<chan::Receiver<Id>, Self::Error>;
} }
/// Public node & device identifier. /// Public node & device identifier.
@ -80,7 +97,11 @@ impl Node {
} }
impl Handle for Node { impl Handle for Node {
fn fetch(&self, id: &Id) -> Result<(), Error> { type Session = ();
type FetchLookup = ();
type Error = Error;
fn fetch(&mut self, id: Id) -> Result<(), Error> {
for line in self.call("fetch", &[id])? { for line in self.call("fetch", &[id])? {
let line = line?; let line = line?;
log::debug!("node: {}", line); log::debug!("node: {}", line);
@ -88,9 +109,9 @@ impl Handle for Node {
Ok(()) Ok(())
} }
fn track_node(&self, id: &NodeId, alias: Option<&str>) -> Result<bool, Error> { fn track_node(&mut self, id: NodeId, alias: Option<String>) -> Result<bool, Error> {
let id = id.to_human(); let id = id.to_human();
let mut line = if let Some(alias) = alias { let mut line = if let Some(alias) = alias.as_deref() {
self.call("track-node", &[id.as_str(), alias]) self.call("track-node", &[id.as_str(), alias])
} else { } else {
self.call("track-node", &[id.as_str()]) self.call("track-node", &[id.as_str()])
@ -111,7 +132,7 @@ impl Handle for Node {
} }
} }
fn track_repo(&self, id: &Id) -> Result<bool, Error> { fn track_repo(&mut self, id: Id) -> Result<bool, Error> {
let mut line = self.call("track-repo", &[id])?; let mut line = self.call("track-repo", &[id])?;
let line = line let line = line
.next() .next()
@ -129,7 +150,7 @@ impl Handle for Node {
} }
} }
fn untrack_node(&self, id: &NodeId) -> Result<bool, Error> { fn untrack_node(&mut self, id: NodeId) -> Result<bool, Error> {
let mut line = self.call("untrack-node", &[id])?; let mut line = self.call("untrack-node", &[id])?;
let line = line.next().ok_or(Error::EmptyResponse { let line = line.next().ok_or(Error::EmptyResponse {
cmd: "untrack-node", cmd: "untrack-node",
@ -147,7 +168,7 @@ impl Handle for Node {
} }
} }
fn untrack_repo(&self, id: &Id) -> Result<bool, Error> { fn untrack_repo(&mut self, id: Id) -> Result<bool, Error> {
let mut line = self.call("untrack-repo", &[id])?; let mut line = self.call("untrack-repo", &[id])?;
let line = line.next().ok_or(Error::EmptyResponse { let line = line.next().ok_or(Error::EmptyResponse {
cmd: "untrack-repo", cmd: "untrack-repo",
@ -165,7 +186,7 @@ impl Handle for Node {
} }
} }
fn announce_refs(&self, id: &Id) -> Result<(), Error> { fn announce_refs(&mut self, id: Id) -> Result<(), Error> {
for line in self.call("announce-refs", &[id])? { for line in self.call("announce-refs", &[id])? {
let line = line?; let line = line?;
log::debug!("node: {}", line); log::debug!("node: {}", line);
@ -173,6 +194,22 @@ impl Handle for Node {
Ok(()) Ok(())
} }
fn routing(&self) -> Result<chan::Receiver<(Id, NodeId)>, Error> {
todo!();
}
fn sessions(&self) -> Result<chan::Receiver<(NodeId, Self::Session)>, Error> {
todo!();
}
fn listening(&self) -> Result<net::SocketAddr, Error> {
todo!();
}
fn inventory(&self) -> Result<chan::Receiver<Id>, Error> {
todo!();
}
fn shutdown(self) -> Result<(), Error> { fn shutdown(self) -> Result<(), Error> {
todo!(); todo!();
} }

View File

@ -197,10 +197,13 @@ pub fn clone<P: AsRef<Path>, G: Signer, S: storage::WriteStorage, H: node::Handl
path: P, path: P,
signer: &G, signer: &G,
storage: &S, storage: &S,
handle: &H, handle: &mut H,
) -> Result<git2::Repository, CloneError> { ) -> Result<git2::Repository, CloneError>
let _ = handle.track_repo(&proj)?; where
let _ = handle.fetch(&proj)?; CloneError: From<H::Error>,
{
let _ = handle.track_repo(proj)?;
let _ = handle.fetch(proj)?;
let _ = fork(proj, signer, storage)?; let _ = fork(proj, signer, storage)?;
let working = checkout(proj, signer.public_key(), path, storage)?; let working = checkout(proj, signer.public_key(), path, storage)?;