Skip to content

Commit cfbcb9b

Browse files
committed
Add support for FD passing to capnp-rpc
Closes: #625
1 parent 118419c commit cfbcb9b

15 files changed

Lines changed: 441 additions & 46 deletions

File tree

capnp-futures/src/io/write_queue.rs

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,7 @@ where
3737
M: AsOutputSegments,
3838
{
3939
Message(M, oneshot::Sender<M>),
40+
MessageMove(M, oneshot::Sender<()>),
4041
Done(Result<(), Error>, oneshot::Sender<()>),
4142
}
4243

@@ -91,6 +92,13 @@ where
9192
writer.flush().await?;
9293
let _ = returner.send(m);
9394
}
95+
Item::MessageMove(m, finisher) => {
96+
let result = crate::io::serialize::write_message(&mut writer, m).await;
97+
in_flight.fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
98+
result?;
99+
writer.flush().await?;
100+
let _ = finisher.send(());
101+
}
94102
Item::Done(r, finisher) => {
95103
let _ = finisher.send(());
96104
return r;
@@ -134,6 +142,21 @@ where
134142
MapErr(oneshot)
135143
}
136144

145+
pub fn send_move(
146+
&mut self,
147+
message: M,
148+
) -> impl Future<Output = Result<(), Error>> + Unpin + 'static {
149+
self.in_flight
150+
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
151+
let (complete, oneshot) = oneshot::channel();
152+
153+
let _ = self
154+
.sender
155+
.unbounded_send(Item::MessageMove(message, complete));
156+
157+
MapErr(oneshot)
158+
}
159+
137160
/// Returns the number of messages queued to be written.
138161
pub fn len(&self) -> usize {
139162
let result = self.in_flight.load(std::sync::atomic::Ordering::SeqCst);

capnp-rpc/src/broken.rs

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
// THE SOFTWARE.
2121

2222
use capnp::any_pointer;
23+
use capnp::fd::BorrowedFd;
2324
use capnp::private::capability::{
2425
ClientHook, ParamsHook, PipelineHook, PipelineOp, RequestHook, ResultsHook,
2526
};
@@ -158,6 +159,10 @@ impl ClientHook for Client {
158159
fn when_resolved(&self) -> Promise<(), Error> {
159160
crate::rpc::default_when_resolved_impl(self)
160161
}
162+
163+
fn get_fd(&self) -> Option<BorrowedFd<'_>> {
164+
None
165+
}
161166
}
162167

163168
pub(crate) fn new_cap(exception: Error) -> Box<dyn ClientHook> {

capnp-rpc/src/lib.rs

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,8 +59,10 @@
5959
//! For a more complete example, see <https://github.com/capnproto/capnproto-rust/tree/master/capnp-rpc/examples/calculator>
6060
6161
use capnp::capability::Promise;
62+
use capnp::fd::FdHooks;
6263
use capnp::private::capability::ClientHook;
6364
use capnp::Error;
65+
use capnp_futures::io::FdReadBuf;
6466
use futures_channel::oneshot;
6567
use futures_util::{FutureExt as _, TryFutureExt as _};
6668
use std::cell::RefCell;
@@ -123,6 +125,10 @@ pub trait OutgoingMessage {
123125
/// Same as `get_body()`, but returns the corresponding reader type.
124126
fn get_body_as_reader(&self) -> ::capnp::Result<::capnp::any_pointer::Reader<'_>>;
125127

128+
fn set_fds(&mut self, fds: FdHooks) {
129+
let _ = fds;
130+
}
131+
126132
/// Sends the message. Returns a promise that resolves once the send has completed.
127133
/// Dropping the returned promise does *not* cancel the send.
128134
fn send(
@@ -149,6 +155,15 @@ pub trait IncomingMessage {
149155
/// The standard RPC implementation interprets it as a Message as defined
150156
/// in `schema/rpc.capnp`.
151157
fn get_body(&self) -> ::capnp::Result<::capnp::any_pointer::Reader<'_>>;
158+
159+
fn get_body_and_attached_fds(
160+
&mut self,
161+
) -> (
162+
::capnp::Result<::capnp::any_pointer::Reader<'_>>,
163+
&mut FdReadBuf,
164+
) {
165+
(self.get_body(), &mut [])
166+
}
152167
}
153168

154169
/// A two-way RPC connection.

capnp-rpc/src/local.rs

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
// THE SOFTWARE.
2020

2121
use capnp::capability::{self, Promise};
22+
use capnp::fd::BorrowedFd;
2223
use capnp::private::capability::{
2324
ClientHook, ParamsHook, PipelineHook, PipelineOp, RequestHook, ResponseHook, ResultsHook,
2425
};
@@ -612,4 +613,8 @@ where
612613
fn when_resolved(&self) -> Promise<(), Error> {
613614
crate::rpc::default_when_resolved_impl(self)
614615
}
616+
617+
fn get_fd(&self) -> Option<BorrowedFd<'_>> {
618+
self.inner.server.get_fd()
619+
}
615620
}

capnp-rpc/src/queued.rs

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121

2222
use capnp::any_pointer;
2323
use capnp::capability::Promise;
24+
use capnp::fd::BorrowedFd;
2425
use capnp::private::capability::{ClientHook, ParamsHook, PipelineHook, PipelineOp, ResultsHook};
2526
use capnp::Error;
2627
use futures_util::{FutureExt as _, TryFutureExt as _};
@@ -352,4 +353,8 @@ impl ClientHook for Client {
352353
fn when_resolved(&self) -> Promise<(), Error> {
353354
crate::rpc::default_when_resolved_impl(self)
354355
}
356+
357+
fn get_fd(&self) -> Option<BorrowedFd<'_>> {
358+
self.inner.redirect.get().and_then(|p| p.get_fd())
359+
}
355360
}

capnp-rpc/src/reconnect.rs

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ use std::marker::PhantomData;
33
use std::rc::Rc;
44

55
use capnp::capability::{FromClientHook, Promise};
6+
use capnp::fd::BorrowedFd;
67
use capnp::private::capability::{ClientHook, RequestHook};
78
use futures_util::TryFutureExt as _;
89

@@ -206,6 +207,13 @@ where
206207
fn when_resolved(&self) -> Promise<(), capnp::Error> {
207208
Promise::ok(())
208209
}
210+
211+
fn get_fd(&self) -> Option<BorrowedFd<'_>> {
212+
// We can’t return `self.get_current().get_fd()` because it might not
213+
// live long enough. This matches the C++ behaviour; see
214+
// <https://github.com/capnproto/capnproto/blob/291d8ee870ab394d32fbdd5d2160add84b65e229/c%2B%2B/src/capnp/reconnect.c%2B%2B#L72-L77>.
215+
None
216+
}
209217
}
210218

211219
struct Request<F, C> {

0 commit comments

Comments
 (0)