protocol/service: rename `Responder::new` to `Responder::oneshot`

The `Responder` is a oneshot channel where it sends one result and the
other side receives that result.

This surfaced that the `AnnounceRefs` command was returning one result
at a time. At the moment, this is only used for a single namespace,
the local user's.
This commit is contained in:
Fintan Halpenny 2026-02-13 14:13:16 +00:00 committed by Lorenz Leutgeb
parent 0d628a45e2
commit 980ed56186
3 changed files with 35 additions and 27 deletions

View File

@ -203,7 +203,7 @@ impl radicle::node::Handle for Handle {
id: RepoId, id: RepoId,
namespaces: impl IntoIterator<Item = PublicKey>, namespaces: impl IntoIterator<Item = PublicKey>,
) -> Result<Seeds, Self::Error> { ) -> Result<Seeds, Self::Error> {
let (responder, receiver) = service::command::Responder::new(); let (responder, receiver) = service::command::Responder::oneshot();
self.command(service::Command::Seeds( self.command(service::Command::Seeds(
id, id,
HashSet::from_iter(namespaces), HashSet::from_iter(namespaces),
@ -213,13 +213,13 @@ impl radicle::node::Handle for Handle {
} }
fn config(&self) -> Result<Config, Self::Error> { fn config(&self) -> Result<Config, Self::Error> {
let (responder, receiver) = service::command::Responder::new(); let (responder, receiver) = service::command::Responder::oneshot();
self.command(service::Command::Config(responder))?; self.command(service::Command::Config(responder))?;
Ok(receiver.recv()??) Ok(receiver.recv()??)
} }
fn listen_addrs(&self) -> Result<Vec<net::SocketAddr>, Self::Error> { fn listen_addrs(&self) -> Result<Vec<net::SocketAddr>, Self::Error> {
let (responder, receiver) = service::command::Responder::new(); let (responder, receiver) = service::command::Responder::oneshot();
self.command(service::Command::ListenAddrs(responder))?; self.command(service::Command::ListenAddrs(responder))?;
Ok(receiver.recv()??) Ok(receiver.recv()??)
} }
@ -230,31 +230,31 @@ impl radicle::node::Handle for Handle {
from: NodeId, from: NodeId,
timeout: time::Duration, timeout: time::Duration,
) -> Result<FetchResult, Error> { ) -> Result<FetchResult, Error> {
let (responder, receiver) = service::command::Responder::new(); let (responder, receiver) = service::command::Responder::oneshot();
self.command(service::Command::Fetch(id, from, timeout, responder))?; self.command(service::Command::Fetch(id, from, timeout, responder))?;
Ok(receiver.recv()??) Ok(receiver.recv()??)
} }
fn follow(&mut self, id: NodeId, alias: Option<Alias>) -> Result<bool, Error> { fn follow(&mut self, id: NodeId, alias: Option<Alias>) -> Result<bool, Error> {
let (responder, receiver) = service::command::Responder::new(); let (responder, receiver) = service::command::Responder::oneshot();
self.command(service::Command::Follow(id, alias, responder))?; self.command(service::Command::Follow(id, alias, responder))?;
Ok(receiver.recv()??) Ok(receiver.recv()??)
} }
fn unfollow(&mut self, id: NodeId) -> Result<bool, Error> { fn unfollow(&mut self, id: NodeId) -> Result<bool, Error> {
let (responder, receiver) = service::command::Responder::new(); let (responder, receiver) = service::command::Responder::oneshot();
self.command(service::Command::Unfollow(id, responder))?; self.command(service::Command::Unfollow(id, responder))?;
Ok(receiver.recv()??) Ok(receiver.recv()??)
} }
fn seed(&mut self, id: RepoId, scope: policy::Scope) -> Result<bool, Error> { fn seed(&mut self, id: RepoId, scope: policy::Scope) -> Result<bool, Error> {
let (responder, receiver) = service::command::Responder::new(); let (responder, receiver) = service::command::Responder::oneshot();
self.command(service::Command::Seed(id, scope, responder))?; self.command(service::Command::Seed(id, scope, responder))?;
Ok(receiver.recv()??) Ok(receiver.recv()??)
} }
fn unseed(&mut self, id: RepoId) -> Result<bool, Error> { fn unseed(&mut self, id: RepoId) -> Result<bool, Error> {
let (responder, receiver) = service::command::Responder::new(); let (responder, receiver) = service::command::Responder::oneshot();
self.command(service::Command::Unseed(id, responder))?; self.command(service::Command::Unseed(id, responder))?;
Ok(receiver.recv()??) Ok(receiver.recv()??)
} }
@ -264,7 +264,7 @@ impl radicle::node::Handle for Handle {
id: RepoId, id: RepoId,
namespaces: impl IntoIterator<Item = PublicKey>, namespaces: impl IntoIterator<Item = PublicKey>,
) -> Result<RefsAt, Error> { ) -> Result<RefsAt, Error> {
let (responder, receiver) = service::command::Responder::new(); let (responder, receiver) = service::command::Responder::oneshot();
self.command(service::Command::AnnounceRefs( self.command(service::Command::AnnounceRefs(
id, id,
HashSet::from_iter(namespaces), HashSet::from_iter(namespaces),
@ -279,7 +279,7 @@ impl radicle::node::Handle for Handle {
} }
fn add_inventory(&mut self, rid: RepoId) -> Result<bool, Error> { fn add_inventory(&mut self, rid: RepoId) -> Result<bool, Error> {
let (responder, receiver) = service::command::Responder::new(); let (responder, receiver) = service::command::Responder::oneshot();
self.command(service::Command::AddInventory(rid, responder))?; self.command(service::Command::AddInventory(rid, responder))?;
Ok(receiver.recv()??) Ok(receiver.recv()??)
} }

View File

@ -874,8 +874,16 @@ where
match self.announce_own_refs(id, doc, namespaces) { match self.announce_own_refs(id, doc, namespaces) {
Ok((refs, _timestamp)) => { Ok((refs, _timestamp)) => {
for r in refs { // TODO(finto): currently the command caller only
resp.ok(r).ok(); // expects one `RefsAt`, this should be fixed in the
// trait, eventually.
if let Some(refs) = refs.first() {
resp.ok(*refs).ok();
} else {
resp.err(command::Error::custom(format!(
"no refs were announced for {id}"
)))
.ok();
} }
} }
Err(err) => { Err(err) => {

View File

@ -24,7 +24,7 @@ pub type Result<T> = std::result::Result<T, Error>;
/// A [`Responder`] returns results after processing a service [`Command`]. /// A [`Responder`] returns results after processing a service [`Command`].
/// ///
/// To construct a [`Responder`], use [`Responder::new`], which also returns its /// To construct a [`Responder`], use [`Responder::oneshot`], which also returns its
/// corresponding [`Receiver`]. /// corresponding [`Receiver`].
/// ///
/// To send results, use either: /// To send results, use either:
@ -38,23 +38,23 @@ pub struct Responder<T> {
impl<T> Responder<T> { impl<T> Responder<T> {
/// Construct a new [`Responder`] and its corresponding [`Receiver`]. /// Construct a new [`Responder`] and its corresponding [`Receiver`].
pub fn new() -> (Self, Receiver<Result<T>>) { pub fn oneshot() -> (Self, Receiver<Result<T>>) {
let (sender, receiver) = crossbeam_channel::bounded(1); let (sender, receiver) = crossbeam_channel::bounded(1);
(Self { channel: sender }, receiver) (Self { channel: sender }, receiver)
} }
/// Send a [`Result`] to the receiver. /// Send a [`Result`] to the receiver.
pub fn send(&self, result: Result<T>) -> std::result::Result<(), SendError<Result<T>>> { pub fn send(self, result: Result<T>) -> std::result::Result<(), SendError<Result<T>>> {
self.channel.send(result) self.channel.send(result)
} }
/// Send a [`Result::Ok`] to the receiver. /// Send a [`Result::Ok`] to the receiver.
pub fn ok(&self, value: T) -> std::result::Result<(), SendError<Result<T>>> { pub fn ok(self, value: T) -> std::result::Result<(), SendError<Result<T>>> {
self.send(Ok(value)) self.send(Ok(value))
} }
/// Send a [`Result::Err`] to the receiver. /// Send a [`Result::Err`] to the receiver.
pub fn err<E>(&self, error: E) -> std::result::Result<(), SendError<Result<T>>> pub fn err<E>(self, error: E) -> std::result::Result<(), SendError<Result<T>>>
where where
E: std::error::Error + Send + Sync + 'static, E: std::error::Error + Send + Sync + 'static,
{ {
@ -108,7 +108,7 @@ impl Command {
rid: RepoId, rid: RepoId,
keys: HashSet<PublicKey>, keys: HashSet<PublicKey>,
) -> (Self, Receiver<Result<RefsAt>>) { ) -> (Self, Receiver<Result<RefsAt>>) {
let (responder, receiver) = Responder::new(); let (responder, receiver) = Responder::oneshot();
(Self::AnnounceRefs(rid, keys, responder), receiver) (Self::AnnounceRefs(rid, keys, responder), receiver)
} }
@ -117,7 +117,7 @@ impl Command {
} }
pub fn add_inventory(rid: RepoId) -> (Self, Receiver<Result<bool>>) { pub fn add_inventory(rid: RepoId) -> (Self, Receiver<Result<bool>>) {
let (responder, receiver) = Responder::new(); let (responder, receiver) = Responder::oneshot();
(Self::AddInventory(rid, responder), receiver) (Self::AddInventory(rid, responder), receiver)
} }
@ -130,17 +130,17 @@ impl Command {
} }
pub fn config() -> (Self, Receiver<Result<Config>>) { pub fn config() -> (Self, Receiver<Result<Config>>) {
let (responder, receiver) = Responder::new(); let (responder, receiver) = Responder::oneshot();
(Self::Config(responder), receiver) (Self::Config(responder), receiver)
} }
pub fn listen_addrs() -> (Self, Receiver<Result<Vec<std::net::SocketAddr>>>) { pub fn listen_addrs() -> (Self, Receiver<Result<Vec<std::net::SocketAddr>>>) {
let (responder, receiver) = Responder::new(); let (responder, receiver) = Responder::oneshot();
(Self::ListenAddrs(responder), receiver) (Self::ListenAddrs(responder), receiver)
} }
pub fn seeds(rid: RepoId, keys: HashSet<PublicKey>) -> (Self, Receiver<Result<Seeds>>) { pub fn seeds(rid: RepoId, keys: HashSet<PublicKey>) -> (Self, Receiver<Result<Seeds>>) {
let (responder, receiver) = Responder::new(); let (responder, receiver) = Responder::oneshot();
(Self::Seeds(rid, keys, responder), receiver) (Self::Seeds(rid, keys, responder), receiver)
} }
@ -149,27 +149,27 @@ impl Command {
node_id: NodeId, node_id: NodeId,
duration: time::Duration, duration: time::Duration,
) -> (Self, Receiver<Result<FetchResult>>) { ) -> (Self, Receiver<Result<FetchResult>>) {
let (responder, receiver) = Responder::new(); let (responder, receiver) = Responder::oneshot();
(Self::Fetch(rid, node_id, duration, responder), receiver) (Self::Fetch(rid, node_id, duration, responder), receiver)
} }
pub fn seed(rid: RepoId, scope: Scope) -> (Self, Receiver<Result<bool>>) { pub fn seed(rid: RepoId, scope: Scope) -> (Self, Receiver<Result<bool>>) {
let (responder, receiver) = Responder::new(); let (responder, receiver) = Responder::oneshot();
(Self::Seed(rid, scope, responder), receiver) (Self::Seed(rid, scope, responder), receiver)
} }
pub fn unseed(rid: RepoId) -> (Self, Receiver<Result<bool>>) { pub fn unseed(rid: RepoId) -> (Self, Receiver<Result<bool>>) {
let (responder, receiver) = Responder::new(); let (responder, receiver) = Responder::oneshot();
(Self::Unseed(rid, responder), receiver) (Self::Unseed(rid, responder), receiver)
} }
pub fn follow(node_id: NodeId, alias: Option<Alias>) -> (Self, Receiver<Result<bool>>) { pub fn follow(node_id: NodeId, alias: Option<Alias>) -> (Self, Receiver<Result<bool>>) {
let (responder, receiver) = Responder::new(); let (responder, receiver) = Responder::oneshot();
(Self::Follow(node_id, alias, responder), receiver) (Self::Follow(node_id, alias, responder), receiver)
} }
pub fn unfollow(node_id: NodeId) -> (Self, Receiver<Result<bool>>) { pub fn unfollow(node_id: NodeId) -> (Self, Receiver<Result<bool>>) {
let (responder, receiver) = Responder::new(); let (responder, receiver) = Responder::oneshot();
(Self::Unfollow(node_id, responder), receiver) (Self::Unfollow(node_id, responder), receiver)
} }