#[cfg(feature = "native")]
pub mod authority;
#[cfg(feature = "replication")]
pub mod authority_admission;
#[cfg(feature = "semantic-analysis")]
pub mod behavior;
#[cfg(feature = "replication")]
pub mod boundary_acceptance;
#[cfg(feature = "replication")]
pub mod collaboration;
pub mod content;
#[cfg(feature = "replication")]
pub mod creation;
#[cfg(feature = "signing")]
pub mod credentials;
#[cfg(feature = "replication")]
pub mod evidence;
#[cfg(feature = "source-transfer")]
pub mod fetch;
#[cfg(feature = "replication")]
pub mod live_replication;
pub mod observation;
#[cfg(feature = "root-attachment")]
pub mod pairing;
#[cfg(any(feature = "native", feature = "replication", feature = "iroh"))]
pub mod publication;
mod reopen;
#[cfg(feature = "replication")]
pub mod replication;
#[cfg(all(feature = "native", feature = "iroh"))]
pub mod replication_rpc;
#[cfg(feature = "signing")]
pub mod request_proof;
#[cfg(feature = "root-attachment")]
pub mod root_attachment;
#[cfg(feature = "replication")]
pub mod thread_control;
#[cfg(feature = "replication")]
pub mod thread_ownership;
pub mod transport;
use api::v2::client::{Client, ClientError, RpcTransport};
pub use api::{
heddle::api::v1alpha2 as contract,
v2::{client::Rpc, rpc},
};
use contract::{DescribeEndpointRequest, DescribeEndpointResponse, EndpointKind, ThreadRef};
pub use reopen::is_reopen_retryable;
use transport::Error;
pub struct Remote<T: RpcTransport<Error = Error>> {
pub api: Client<T>,
pub description: DescribeEndpointResponse,
}
impl<T: RpcTransport<Error = Error>> Remote<T> {
pub async fn discover(
transport: T,
endpoint_key: [u8; 32],
kind: EndpointKind,
) -> Result<Self, ClientError<Error>> {
let bytes = transport
.unary(
rpc::EndpointServiceDescribeEndpoint::METHOD,
prost::Message::encode_to_vec(&DescribeEndpointRequest {
understood_packages: vec!["heddle.api.v1alpha2".into()],
}),
)
.await
.map_err(ClientError::Transport)?;
let description: DescribeEndpointResponse = prost::Message::decode(bytes.as_slice())?;
if description
.endpoint
.as_ref()
.is_none_or(|source| source.public_key != endpoint_key || source.kind != kind as i32)
|| !description
.supported_packages
.iter()
.any(|p| p == "heddle.api.v1alpha2")
{
return Err(ClientError::Transport(Error::Protocol(
"endpoint identity/package mismatch",
)));
}
let api = Client::new(transport, description.implemented_methods.clone());
Ok(Self { api, description })
}
pub async fn observe<M>(
&self,
mut request: M::Request,
resume: Option<observation::Resume>,
) -> Result<observation::Observation<T::Reader, M::Response>, observation::Error>
where
M: api::v2::client::ServerStreamingRpc,
M::Request: observation::ObservationRequest,
M::Response: observation::ObservedEvent,
{
use crate::reopen::ReopenRetryable as _;
use observation::ObservationRequest as _;
let budget = observation::budget(&self.description)?;
let options = request.options_mut();
options.budget = Some(budget);
options.after_cursor.clear();
let mut query = b"heddle-observation-query-v2\0".to_vec();
query.extend_from_slice(M::METHOD.path.as_bytes());
query.push(0);
query.extend_from_slice(&prost::Message::encode_to_vec(&request));
observation::validate_resume(&resume, &self.description, &query)?;
if let Some(resume) = &resume {
request.options_mut().after_cursor = resume.cursor.clone();
}
let mut attempt = 0;
loop {
let messages = match self.api.observe::<M>(&request).await {
Ok(messages) => messages,
Err(error)
if reopen::client_error_is_reopen_retryable(&error)
&& attempt + 1 < reopen::ATTEMPTS =>
{
attempt += 1;
reopen::backoff(attempt).await;
continue;
}
Err(error) => return Err(error.into()),
};
let mut observation = observation::Observation::new(
messages,
&self.description,
budget,
resume.clone(),
query.clone(),
)?;
match observation.consume_open().await {
Ok(()) => return Ok(observation),
Err(error) if error.is_reopen_retryable() && attempt + 1 < reopen::ATTEMPTS => {
attempt += 1;
reopen::backoff(attempt).await;
}
Err(error) => {
observation.prime_error(error);
return Ok(observation);
}
}
}
}
pub async fn observe_analysis(
&self,
request: contract::ObserveAnalysisRequest,
resume: Option<observation::Resume>,
) -> Result<observation::AnalysisObservation<T::Reader>, observation::Error> {
self.observe::<rpc::AnalysisServiceObserveAnalysis>(request, resume)
.await
}
pub fn thread(&self, thread: ThreadRef) -> Thread<'_, T> {
Thread {
remote: self,
reference: thread,
}
}
}
pub struct Thread<'a, T: RpcTransport<Error = Error>> {
remote: &'a Remote<T>,
pub reference: ThreadRef,
}
impl<T: RpcTransport<Error = Error>> Thread<'_, T> {
#[cfg(feature = "replication")]
pub async fn revise_intent(
&self,
command: &thread_control::PreparedControl,
) -> Result<contract::ThreadMutationResponse, ClientError<Error>> {
let request = command.revise_intent().map_err(ClientError::Transport)?;
if request.thread.as_ref() != Some(&self.reference) {
return Err(ClientError::Transport(Error::Protocol(
"prepared command belongs to another Thread",
)));
}
self.remote
.api
.call::<rpc::ThreadServiceReviseIntent>(&request)
.await
}
pub async fn observe(
&self,
sections: &[contract::ThreadSection],
mode: contract::ObservationMode,
resume: Option<observation::Resume>,
) -> Result<observation::ThreadObservation<T::Reader>, observation::Error> {
self.remote
.observe::<rpc::ThreadServiceObserveThread>(
contract::ObserveThreadRequest {
thread: Some(self.reference.clone()),
sections: sections.iter().map(|section| *section as i32).collect(),
observe: Some(contract::ObserveOptions {
mode: mode as i32,
..Default::default()
}),
..Default::default()
},
resume,
)
.await
}
}