use capnp::private::capability::{ClientHook, ParamsHook, ResultsHook};
use capnp::traits::ImbueMut;
use crate::error::{Result, RpcError};
use crate::io::{AsyncStream, read_message, write_raw};
use crate::rpc_capnp;
use crate::tunnelrpc_capnp;
pub const CLOUDFLARED_SERVER_INTERFACE_ID: u64 = 0xf548_cef9_dea2_a4a1;
pub const SESSION_MANAGER_INTERFACE_ID: u64 = 0x8394_45a5_9fb0_1686;
pub const CONFIGURATION_MANAGER_INTERFACE_ID: u64 = 0xb48e_dfbd_aa25_db04;
#[derive(Debug, Clone, Default)]
pub struct UpdateConfigurationResponse {
pub latest_applied_version: i32,
pub error: String,
}
#[derive(Debug, Clone, Default)]
pub struct RegisterUdpSessionResponse {
pub error: String,
pub spans: Vec<u8>,
}
pub trait CloudflaredHandler: Send + Sync {
fn update_configuration(
&self,
version: i32,
configuration: &[u8],
) -> UpdateConfigurationResponse;
fn register_udp_session(
&self,
_session_identifier: &[u8; 16],
_destination_ip: &[u8],
_destination_port: u16,
) -> RegisterUdpSessionResponse {
RegisterUdpSessionResponse {
error: "UDP sessions are not supported".into(),
spans: Vec::new(),
}
}
fn unregister_udp_session(&self, _session_identifier: &[u8; 16], _message: &str) {}
}
pub async fn serve_cloudflared<S, H>(stream: &mut S, handler: &H) -> Result<()>
where
S: AsyncStream + Unpin,
H: CloudflaredHandler,
{
loop {
let reply: Option<Vec<u8>> = {
let reader = match read_message(stream).await {
Ok(reader) => reader,
Err(RpcError::Eof) => return Ok(()),
Err(e) => return Err(e),
};
let root = reader.get_root::<rpc_capnp::message::Reader>()?;
match root.reborrow().which()? {
rpc_capnp::message::Bootstrap(b) => {
let question = b?.get_question_id();
tracing::debug!(question, "edge bootstrapped the RPC stream");
Some(build_bootstrap_return(question)?)
}
rpc_capnp::message::Call(c) => {
let call = c?;
let question = call.get_question_id();
tracing::debug!(
question,
interface_id = call.get_interface_id(),
method_id = call.get_method_id(),
"edge called an RPC method"
);
let reply = match classify(call.get_interface_id(), call.get_method_id()) {
Method::UpdateConfiguration => {
let parameters = call.reborrow().get_params()?.get_content().get_as::<
tunnelrpc_capnp::configuration_manager::update_configuration_params::Reader<'_>,
>()?;
let response = handler.update_configuration(
parameters.get_version(),
parameters.get_config()?,
);
build_update_configuration_return(question, &response)?
}
Method::RegisterUdpSession => {
let parameters = call.reborrow().get_params()?.get_content().get_as::<
tunnelrpc_capnp::session_manager::register_udp_session_params::Reader<'_>,
>()?;
let response = handler.register_udp_session(
&session_identifier_bytes(parameters.get_session_id()?),
parameters.get_dst_ip()?,
parameters.get_dst_port(),
);
build_register_udp_session_return(question, &response)?
}
Method::UnregisterUdpSession => {
build_unregister_udp_session_return(question)?
}
Method::Unknown => {
tracing::debug!(
interface_id = call.get_interface_id(),
method_id = call.get_method_id(),
"edge called an unknown RPC method"
);
crate::rpc::build_exception(question, "unimplemented")?
}
};
Some(reply)
}
rpc_capnp::message::Finish(_) | rpc_capnp::message::Release(_) => None,
_ => None,
}
};
if let Some(reply) = reply {
write_raw(stream, &reply).await?;
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Method {
UpdateConfiguration,
RegisterUdpSession,
UnregisterUdpSession,
Unknown,
}
fn classify(interface_identifier: u64, method_identifier: u16) -> Method {
match (interface_identifier, method_identifier) {
(CONFIGURATION_MANAGER_INTERFACE_ID, 0) | (CLOUDFLARED_SERVER_INTERFACE_ID, 2) => {
Method::UpdateConfiguration
}
(SESSION_MANAGER_INTERFACE_ID, 0) | (CLOUDFLARED_SERVER_INTERFACE_ID, 0) => {
Method::RegisterUdpSession
}
(SESSION_MANAGER_INTERFACE_ID, 1) | (CLOUDFLARED_SERVER_INTERFACE_ID, 1) => {
Method::UnregisterUdpSession
}
_ => Method::Unknown,
}
}
fn session_identifier_bytes(data: &[u8]) -> [u8; 16] {
let mut identifier = [0u8; 16];
let n = data.len().min(16);
identifier[..n].copy_from_slice(&data[..n]);
identifier
}
fn build_bootstrap_return(question: u32) -> Result<Vec<u8>> {
let mut caps: capnp::private::layout::CapTable = Vec::new();
let mut message = capnp::message::Builder::new_default();
let root = message.init_root::<rpc_capnp::message::Builder>();
let mut ret = root.init_return();
ret.set_answer_id(question);
let mut results = ret.init_results();
let mut payload = results.reborrow();
let mut content = payload.reborrow().init_content();
content.imbue_mut(&mut caps);
content.set_as_capability(Box::new(DummyHook));
let mut ctab = payload.reborrow().init_cap_table(1);
ctab.reborrow().get(0).set_sender_hosted(0);
Ok(crate::io::serialize_message(&message))
}
fn build_update_configuration_return(
question: u32,
response: &UpdateConfigurationResponse,
) -> Result<Vec<u8>> {
let mut message = capnp::message::Builder::new_default();
let root = message.init_root::<rpc_capnp::message::Builder>();
let mut ret = root.init_return();
ret.set_answer_id(question);
let mut results = ret.init_results();
let mut payload = results.reborrow();
{
let content = payload.reborrow().init_content();
let mut results_reader =
content.init_as::<tunnelrpc_capnp::configuration_manager::update_configuration_results::Builder>();
let mut response_builder = results_reader.reborrow().init_result();
response_builder.set_latest_applied_version(response.latest_applied_version);
response_builder.set_err(&response.error);
}
payload.reborrow().init_cap_table(0);
Ok(crate::io::serialize_message(&message))
}
fn build_register_udp_session_return(
question: u32,
response: &RegisterUdpSessionResponse,
) -> Result<Vec<u8>> {
let mut message = capnp::message::Builder::new_default();
let root = message.init_root::<rpc_capnp::message::Builder>();
let mut ret = root.init_return();
ret.set_answer_id(question);
let mut results = ret.init_results();
let mut payload = results.reborrow();
{
let content = payload.reborrow().init_content();
let mut results_reader =
content
.init_as::<tunnelrpc_capnp::session_manager::register_udp_session_results::Builder>(
);
let mut response_builder = results_reader.reborrow().init_result();
response_builder.set_err(&response.error);
response_builder.set_spans(&response.spans);
}
payload.reborrow().init_cap_table(0);
Ok(crate::io::serialize_message(&message))
}
fn build_unregister_udp_session_return(question: u32) -> Result<Vec<u8>> {
let mut message = capnp::message::Builder::new_default();
let root = message.init_root::<rpc_capnp::message::Builder>();
let mut ret = root.init_return();
ret.set_answer_id(question);
let mut results = ret.init_results();
let mut payload = results.reborrow();
payload
.reborrow()
.init_content()
.init_as::<tunnelrpc_capnp::session_manager::unregister_udp_session_results::Builder>();
payload.reborrow().init_cap_table(0);
Ok(crate::io::serialize_message(&message))
}
struct DummyHook;
impl ClientHook for DummyHook {
fn add_ref(&self) -> Box<dyn ClientHook> {
Box::new(DummyHook)
}
fn new_call(
&self,
_interface_identifier: u64,
_method_identifier: u16,
_size_hint: Option<capnp::MessageSize>,
) -> capnp::capability::Request<capnp::any_pointer::Owned, capnp::any_pointer::Owned> {
unreachable!("dummy hook is never called")
}
fn call(
&self,
_interface_identifier: u64,
_method_identifier: u16,
_params: Box<dyn ParamsHook>,
_results: Box<dyn ResultsHook>,
) -> capnp::capability::Promise<(), capnp::Error> {
unreachable!("dummy hook is never called")
}
fn get_brand(&self) -> usize {
0
}
fn get_ptr(&self) -> usize {
0
}
fn get_resolved(&self) -> Option<Box<dyn ClientHook>> {
None
}
fn when_more_resolved(
&self,
) -> Option<capnp::capability::Promise<Box<dyn ClientHook>, capnp::Error>> {
None
}
fn when_resolved(&self) -> capnp::capability::Promise<(), capnp::Error> {
capnp::capability::Promise::ok(())
}
}