From c8a24d558416661fc3b39e9dffd323b71a23f224 Mon Sep 17 00:00:00 2001 From: Fintan Halpenny Date: Tue, 23 Apr 2024 17:40:53 +0100 Subject: [PATCH] node: report upload-pack progress Emit the `UploadPack` events from the `upload_pack` method. A wrapper type, `Reporter`, is introduced for implementing a `Write` instance that will parse the given bytes as `gix_protocol::RemoteProgress`, if possible, and emit it via the `Emitter`, and then write the bytes to the underlying writer. This allows any subscribers to see the progress of an upload-pack, i.e. `Compressing objects`, `Counting object`, `Enumerating objects`. Signed-off-by: Fintan Halpenny X-Clacks-Overhead: GNU Terry Pratchett --- Cargo.lock | 1 + radicle-node/Cargo.toml | 1 + radicle-node/src/service.rs | 4 ++ radicle-node/src/wire/protocol.rs | 5 +- radicle-node/src/worker.rs | 25 ++++++--- radicle-node/src/worker/upload_pack.rs | 75 ++++++++++++++++++++++++-- 6 files changed, 99 insertions(+), 12 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index f1e9c3dd..9343e5cb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2537,6 +2537,7 @@ dependencies = [ "crossbeam-channel", "cyphernet", "fastrand", + "gix-protocol", "io-reactor", "lexopt", "libc", diff --git a/radicle-node/Cargo.toml b/radicle-node/Cargo.toml index b873f4f8..cda98b6d 100644 --- a/radicle-node/Cargo.toml +++ b/radicle-node/Cargo.toml @@ -21,6 +21,7 @@ colored = { version = "2.1.0" } crossbeam-channel = { version = "0.5.6" } cyphernet = { version = "0.5.0", features = ["tor", "dns", "ed25519", "p2p-ed25519"] } fastrand = { version = "2.0.0" } +gix-protocol = { version = "0.41.1", features = ["blocking-client"] } io-reactor = { version = "0.5.1", features = ["popol"] } lexopt = { version = "0.3.0" } libc = { version = "0.2.137" } diff --git a/radicle-node/src/service.rs b/radicle-node/src/service.rs index bb1f36e6..5ca26b25 100644 --- a/radicle-node/src/service.rs +++ b/radicle-node/src/service.rs @@ -425,6 +425,10 @@ where pub fn local_time(&self) -> LocalTime { self.clock } + + pub fn emitter(&self) -> Emitter { + self.emitter.clone() + } } impl Service diff --git a/radicle-node/src/wire/protocol.rs b/radicle-node/src/wire/protocol.rs index 442e3e5b..b9575fbb 100644 --- a/radicle-node/src/wire/protocol.rs +++ b/radicle-node/src/wire/protocol.rs @@ -721,7 +721,10 @@ where }; let task = Task { - fetch: FetchRequest::Responder { remote: *nid }, + fetch: FetchRequest::Responder { + remote: *nid, + emitter: self.service.emitter(), + }, stream, channels, }; diff --git a/radicle-node/src/worker.rs b/radicle-node/src/worker.rs index fce17e94..f2340fa8 100644 --- a/radicle-node/src/worker.rs +++ b/radicle-node/src/worker.rs @@ -11,14 +11,14 @@ use std::path::PathBuf; use crossbeam_channel as chan; use radicle::identity::RepoId; -use radicle::node::notifications; +use radicle::node::{notifications, Event}; use radicle::prelude::NodeId; use radicle::storage::refs::RefsAt; use radicle::storage::{ReadRepository, ReadStorage}; use radicle::{cob, crypto, Storage}; use radicle_fetch::FetchLimit; -use crate::runtime::{thread, Handle}; +use crate::runtime::{thread, Emitter, Handle}; use crate::service::policy; use crate::service::policy::Policy; use crate::wire::StreamId; @@ -113,13 +113,15 @@ pub enum FetchRequest { Responder { /// Remote peer we are interacting with. remote: NodeId, + /// Reporter for upload-pack progress. + emitter: Emitter, }, } impl FetchRequest { pub fn remote(&self) -> NodeId { match self { - Self::Initiator { remote, .. } | Self::Responder { remote } => *remote, + Self::Initiator { remote, .. } | Self::Responder { remote, .. } => *remote, } } } @@ -233,7 +235,7 @@ impl Worker { let result = self.fetch(rid, remote, refs_at, channels, notifs); FetchResult::Initiator { rid, result } } - FetchRequest::Responder { remote } => { + FetchRequest::Responder { remote, emitter } => { log::debug!(target: "worker", "Worker processing incoming fetch for {remote} on stream {stream}.."); let (mut stream_r, stream_w) = channels.split(); @@ -255,10 +257,17 @@ impl Worker { }; } - let result = - upload_pack::upload_pack(&self.nid, &self.storage, &header, stream_r, stream_w) - .map(|_| ()) - .map_err(|e| e.into()); + let result = upload_pack::upload_pack( + &self.nid, + remote, + &self.storage, + &emitter, + &header, + stream_r, + stream_w, + ) + .map(|_| ()) + .map_err(|e| e.into()); log::debug!(target: "worker", "Upload process on stream {stream} exited with result {result:?}"); FetchResult::Responder { diff --git a/radicle-node/src/worker/upload_pack.rs b/radicle-node/src/worker/upload_pack.rs index c977a7fd..731b278d 100644 --- a/radicle-node/src/worker/upload_pack.rs +++ b/radicle-node/src/worker/upload_pack.rs @@ -1,8 +1,13 @@ use std::io; use std::io::Write; use std::process::{Command, ExitStatus, Stdio}; +use std::time::Instant; -use radicle::node::NodeId; +use gix_protocol::transport::bstr::ByteSlice; +use radicle::identity::RepoId; +use radicle::node::events; +use radicle::node::events::Emitter; +use radicle::node::{Event, NodeId}; use radicle::storage::git::paths; use radicle::Storage; @@ -16,15 +21,18 @@ use crate::runtime::thread; /// send the EOF file message. pub fn upload_pack( nid: &NodeId, + remote: NodeId, storage: &Storage, + emitter: &Emitter, header: &pktline::GitRequest, mut recv: R, - mut send: W, + send: W, ) -> io::Result where R: io::Read + Send, W: io::Write + Send, { + let timer = Instant::now(); let protocol_version = header .extra .iter() @@ -75,14 +83,21 @@ where let mut stdin = child.stdin.take().unwrap(); let mut stdout = io::BufReader::new(child.stdout.take().unwrap()); + let mut reporter = Reporter { + rid: header.repo, + remote, + emitter: emitter.clone(), + send, + }; thread::scope(|s| { thread::spawn_scoped(nid, "upload-pack", s, || { // N.b. we indefinitely copy stdout to the sender, // i.e. there's no need for a loop. - match io::copy(&mut stdout, &mut send) { + match io::copy(&mut stdout, &mut reporter) { Ok(_) => {} Err(e) => { log::error!(target: "worker", "Worker channel disconnected for {}; aborting: {e}", header.repo); + emitter.emit(events::UploadPack::error(header.repo, remote, e).into()); } } }); @@ -104,6 +119,7 @@ where } Err(e) => { log::error!(target: "worker", "Error on upload-pack channel read for {}: {e}", header.repo); + emitter.emit(events::UploadPack::error(header.repo, remote, e).into()); break; } } @@ -124,9 +140,62 @@ where })?; let status = child.wait()?; + emitter.emit(events::UploadPack::done(header.repo, remote, status).into()); + log::debug!(target: "worker", "Upload pack finished ({}ms)", timer.elapsed().as_millis()); Ok(status) } +/// A combination of the upload-pack sender with an [`Emitter`] for reporting +/// the progress events to subscribers. +struct Reporter { + rid: RepoId, + remote: NodeId, + emitter: Emitter, + send: W, +} + +impl Reporter { + fn emit(&self, buf: &[u8]) { + if let Some(progress) = Self::as_upload_pack_progress(buf) { + log::trace!(target: "worker", "upload-pack progress: {progress}"); + self.emitter + .emit(events::UploadPack::write(self.rid, self.remote, progress).into()); + } + } + + fn as_upload_pack_progress(buf: &[u8]) -> Option { + use events::upload_pack::Progress::*; + let gix_protocol::RemoteProgress { + action, step, max, .. + } = gix_protocol::RemoteProgress::from_bytes(buf)?; + if action.contains_str("Counting objects") { + step.and_then(|processed| max.map(|total| Counting { processed, total })) + } else if action.contains_str("Compressing objects") { + step.and_then(|processed| max.map(|total| Compressing { processed, total })) + } else if action.contains_str("Enumerating objects") { + max.map(|total| Enumerating { total }) + } else { + None + } + } +} + +impl io::Write for Reporter +where + W: io::Write, +{ + fn write(&mut self, buf: &[u8]) -> io::Result { + let n = self.send.write(buf)?; + self.send.flush()?; + self.emit(buf); + Ok(n) + } + + fn flush(&mut self) -> io::Result<()> { + self.send.flush() + } +} + pub(super) mod pktline { use std::io; use std::io::Read;