331 lines
11 KiB
Rust
331 lines
11 KiB
Rust
mod handle;
|
||
|
||
use std::io::{BufRead, BufReader};
|
||
use std::os::unix::net::UnixListener;
|
||
use std::path::PathBuf;
|
||
use std::sync::{Arc, Mutex};
|
||
use std::{fs, io, net, thread, time};
|
||
|
||
use crossbeam_channel as chan;
|
||
use cyphernet::Ecdh;
|
||
use netservices::resource::NetAccept;
|
||
use reactor::poller::popol;
|
||
use reactor::Reactor;
|
||
use thiserror::Error;
|
||
|
||
use radicle::git;
|
||
use radicle::node::Handle as _;
|
||
use radicle::node::{ADDRESS_DB_FILE, ROUTING_DB_FILE, TRACKING_DB_FILE};
|
||
use radicle::profile::Home;
|
||
use radicle::Storage;
|
||
|
||
use crate::address;
|
||
use crate::control;
|
||
use crate::crypto::Signer;
|
||
use crate::node::{routing, NodeId};
|
||
use crate::service::{tracking, Event};
|
||
use crate::wire;
|
||
use crate::wire::Wire;
|
||
use crate::worker;
|
||
use crate::{service, LocalTime};
|
||
|
||
pub use handle::Error as HandleError;
|
||
pub use handle::Handle;
|
||
|
||
/// A client error.
|
||
#[derive(Error, Debug)]
|
||
pub enum Error {
|
||
/// A routing database error.
|
||
#[error("routing database error: {0}")]
|
||
Routing(#[from] routing::Error),
|
||
/// An address database error.
|
||
#[error("address database error: {0}")]
|
||
Addresses(#[from] address::Error),
|
||
/// A tracking database error.
|
||
#[error("tracking database error: {0}")]
|
||
Tracking(#[from] tracking::Error),
|
||
/// An I/O error.
|
||
#[error("i/o error: {0}")]
|
||
Io(#[from] io::Error),
|
||
/// A control socket error.
|
||
#[error("control socket error: {0}")]
|
||
Control(#[from] control::Error),
|
||
/// Another node is already running.
|
||
#[error(
|
||
"another node appears to be running; \
|
||
if this isn't the case, delete the socket file at '{0}' \
|
||
and restart the node"
|
||
)]
|
||
AlreadyRunning(PathBuf),
|
||
/// A git version error.
|
||
#[error("git version error: {0}")]
|
||
GitVersion(#[from] git::VersionError),
|
||
}
|
||
|
||
/// Publishes events to subscribers.
|
||
#[derive(Debug, Clone)]
|
||
pub struct Emitter<T> {
|
||
subscribers: Arc<Mutex<Vec<chan::Sender<T>>>>,
|
||
}
|
||
|
||
impl<T> Default for Emitter<T> {
|
||
fn default() -> Emitter<T> {
|
||
Emitter {
|
||
subscribers: Default::default(),
|
||
}
|
||
}
|
||
}
|
||
|
||
impl<T: Clone> Emitter<T> {
|
||
/// Emit event to subscribers and drop those who can't receive it.
|
||
pub(crate) fn emit(&self, event: T) {
|
||
self.subscribers
|
||
.lock()
|
||
.unwrap()
|
||
.retain(|s| s.try_send(event.clone()).is_ok());
|
||
}
|
||
|
||
/// Subscribe to events stream.
|
||
pub fn subscribe(&self) -> chan::Receiver<T> {
|
||
let (sender, receiver) = chan::unbounded();
|
||
let mut subs = self.subscribers.lock().unwrap();
|
||
subs.push(sender);
|
||
|
||
receiver
|
||
}
|
||
}
|
||
|
||
/// Holds join handles to the client threads, as well as a client handle.
|
||
pub struct Runtime {
|
||
pub id: NodeId,
|
||
pub home: Home,
|
||
pub control: UnixListener,
|
||
pub handle: Handle,
|
||
pub storage: Storage,
|
||
pub reactor: Reactor<wire::Control>,
|
||
pub daemon: net::SocketAddr,
|
||
pub pool: worker::Pool,
|
||
pub local_addrs: Vec<net::SocketAddr>,
|
||
pub signals: chan::Receiver<()>,
|
||
}
|
||
|
||
impl Runtime {
|
||
/// Initialize the runtime.
|
||
///
|
||
/// This function spawns threads.
|
||
pub fn init<G: Signer + Ecdh + 'static>(
|
||
home: Home,
|
||
config: service::Config,
|
||
listen: Vec<net::SocketAddr>,
|
||
proxy: net::SocketAddr,
|
||
daemon: net::SocketAddr,
|
||
signals: chan::Receiver<()>,
|
||
signer: G,
|
||
) -> Result<Runtime, Error>
|
||
where
|
||
G: Ecdh<Pk = NodeId> + Clone,
|
||
{
|
||
let id = *signer.public_key();
|
||
let node_dir = home.node();
|
||
let network = config.network;
|
||
let rng = fastrand::Rng::new();
|
||
let clock = LocalTime::now();
|
||
let storage = Storage::open(home.storage())?;
|
||
let address_db = node_dir.join(ADDRESS_DB_FILE);
|
||
let routing_db = node_dir.join(ROUTING_DB_FILE);
|
||
let tracking_db = node_dir.join(TRACKING_DB_FILE);
|
||
|
||
log::info!(target: "node", "Opening address book {}..", address_db.display());
|
||
let addresses = address::Book::open(address_db)?;
|
||
|
||
log::info!(target: "node", "Opening routing table {}..", routing_db.display());
|
||
let routing = routing::Table::open(routing_db)?;
|
||
|
||
log::info!(target: "node", "Opening tracking policy table {}..", tracking_db.display());
|
||
let tracking = tracking::Store::open(tracking_db)?;
|
||
let tracking = tracking::Config::new(config.policy, config.scope, tracking);
|
||
|
||
log::info!(target: "node", "Default tracking policy set to '{}'", &config.policy);
|
||
log::info!(target: "node", "Initializing service ({:?})..", network);
|
||
let emitter: Emitter<Event> = Default::default();
|
||
let service = service::Service::new(
|
||
config,
|
||
clock,
|
||
routing,
|
||
storage.clone(),
|
||
addresses,
|
||
tracking,
|
||
signer.clone(),
|
||
rng,
|
||
emitter.clone(),
|
||
);
|
||
|
||
let (worker_send, worker_recv) = chan::unbounded::<worker::Task>();
|
||
let mut wire = Wire::new(service, worker_send, signer, proxy, clock);
|
||
let mut local_addrs = Vec::new();
|
||
|
||
for addr in listen {
|
||
let listener = NetAccept::bind(&addr)?;
|
||
let local_addr = listener.local_addr();
|
||
|
||
local_addrs.push(local_addr);
|
||
wire.listen(listener);
|
||
|
||
log::info!(target: "node", "Listening on {local_addr}..");
|
||
}
|
||
let reactor = Reactor::named(wire, popol::Poller::new(), id.to_human())?;
|
||
let handle = Handle::new(home.clone(), reactor.controller(), emitter);
|
||
let atomic = git::version()? >= git::VERSION_REQUIRED;
|
||
|
||
if !atomic {
|
||
log::warn!(
|
||
target: "node",
|
||
"Disabling atomic fetches; git version >= {} required", git::VERSION_REQUIRED
|
||
);
|
||
}
|
||
|
||
let pool = worker::Pool::with(
|
||
id,
|
||
worker_recv,
|
||
handle.clone(),
|
||
worker::Config {
|
||
capacity: 8,
|
||
name: id.to_human(),
|
||
timeout: time::Duration::from_secs(9),
|
||
storage: storage.clone(),
|
||
daemon,
|
||
atomic,
|
||
},
|
||
);
|
||
let control = match UnixListener::bind(home.socket()) {
|
||
Ok(sock) => sock,
|
||
Err(err) if err.kind() == io::ErrorKind::AddrInUse => {
|
||
return Err(Error::AlreadyRunning(home.socket()));
|
||
}
|
||
Err(err) => {
|
||
return Err(err.into());
|
||
}
|
||
};
|
||
|
||
Ok(Runtime {
|
||
id,
|
||
home,
|
||
control,
|
||
storage,
|
||
reactor,
|
||
daemon,
|
||
handle,
|
||
pool,
|
||
signals,
|
||
local_addrs,
|
||
})
|
||
}
|
||
|
||
pub fn run(self) -> Result<(), Error> {
|
||
let home = self.home;
|
||
|
||
log::info!(target: "node", "Running node {} in {}..", self.id, home.path().display());
|
||
log::info!(target: "node", "Binding control socket {}..", home.socket().display());
|
||
|
||
let control = thread::Builder::new().name(self.id.to_human()).spawn({
|
||
let handle = self.handle.clone();
|
||
move || control::listen(self.control, handle)
|
||
})?;
|
||
let _signals = thread::Builder::new()
|
||
.name(self.id.to_human())
|
||
.spawn(move || {
|
||
if let Ok(()) = self.signals.recv() {
|
||
log::info!(target: "node", "Termination signal received; shutting down..");
|
||
self.handle.shutdown().ok();
|
||
}
|
||
})?;
|
||
|
||
log::info!(target: "node", "Spawning git daemon at {}..", self.storage.path().display());
|
||
|
||
let mut daemon = daemon::spawn(self.storage.path(), self.daemon)?;
|
||
thread::Builder::new().name(self.id.to_human()).spawn({
|
||
let stderr = daemon.stderr.take().unwrap();
|
||
|| {
|
||
for line in BufReader::new(stderr).lines().flatten() {
|
||
if line.starts_with("fatal") {
|
||
log::error!(target: "daemon", "{line}");
|
||
} else {
|
||
log::debug!(target: "daemon", "{line}");
|
||
}
|
||
}
|
||
}
|
||
})?;
|
||
|
||
self.pool.run().unwrap();
|
||
self.reactor.join().unwrap();
|
||
|
||
daemon::kill(&daemon).ok(); // Ignore error if daemon has already exited, for whatever reason.
|
||
daemon.wait()?;
|
||
|
||
// If the socket file was deleted by some other process, for whatever reason,
|
||
// the control thread will not be able to join.
|
||
if fs::remove_file(home.socket()).is_ok() {
|
||
control.join().unwrap()?;
|
||
}
|
||
log::debug!(target: "node", "Node shutdown completed for {}", self.id);
|
||
|
||
Ok(())
|
||
}
|
||
}
|
||
|
||
pub mod daemon {
|
||
use std::path::Path;
|
||
use std::process::{Child, Command, Stdio};
|
||
use std::{env, io, net};
|
||
|
||
/// Kill the daemon process.
|
||
pub fn kill(child: &Child) -> io::Result<()> {
|
||
// SAFETY: We use `libc::kill` because `Child::kill` always sends a `SIGKILL` and that doesn't
|
||
// work for us. We need to send a `SIGTERM` to fully reap the child process. This is because
|
||
// `git-daemon` spawns its own children, and isn't able to reap them if it receives
|
||
// a `SIGKILL`.
|
||
let result = unsafe { libc::kill(child.id() as libc::c_int, libc::SIGTERM) };
|
||
match result {
|
||
0 => Ok(()),
|
||
_ => Err(io::Error::last_os_error()),
|
||
}
|
||
}
|
||
|
||
/// Spawn the daemon process.
|
||
pub fn spawn(storage: &Path, addr: net::SocketAddr) -> io::Result<Child> {
|
||
let storage = storage.canonicalize()?;
|
||
let listen = format!("--listen={}", addr.ip());
|
||
let port = format!("--port={}", addr.port());
|
||
let child = Command::new("git")
|
||
.env_clear()
|
||
.envs(env::vars().filter(|(k, _)| k == "PATH" || k.starts_with("GIT")))
|
||
.envs(radicle::git::env::GIT_DEFAULT_CONFIG)
|
||
.env("GIT_PROTOCOL", "version=2")
|
||
.current_dir(storage)
|
||
.arg("daemon")
|
||
// Make all git directories available.
|
||
.arg("--export-all")
|
||
.arg("--reuseaddr")
|
||
.arg("--max-connections=32")
|
||
.arg("--informative-errors")
|
||
.arg("--verbose")
|
||
// The git "root". Should be our storage path.
|
||
.arg("--base-path=.")
|
||
// Timeout (in seconds) between the moment the connection is established
|
||
// and the client request is received (typically a rather low value,
|
||
// since that should be basically immediate).
|
||
.arg("--init-timeout=3")
|
||
// Timeout (in seconds) for specific client sub-requests.
|
||
// This includes the time it takes for the server to process the sub-request
|
||
// and the time spent waiting for the next client’s request.
|
||
.arg("--timeout=9")
|
||
.arg("--log-destination=stderr")
|
||
.arg(listen)
|
||
.arg(port)
|
||
.stderr(Stdio::piped())
|
||
.spawn()?;
|
||
|
||
Ok(child)
|
||
}
|
||
}
|