1#[cfg(feature = "native")]
5pub mod authority;
6#[cfg(feature = "replication")]
7pub mod authority_admission;
8#[cfg(feature = "semantic-analysis")]
9pub mod behavior;
10#[cfg(feature = "replication")]
11pub mod boundary_acceptance;
12#[cfg(feature = "replication")]
13pub mod collaboration;
14pub mod content;
15#[cfg(feature = "replication")]
16pub mod creation;
17#[cfg(feature = "signing")]
18pub mod credentials;
19#[cfg(feature = "replication")]
20pub mod evidence;
21#[cfg(feature = "source-transfer")]
22pub mod fetch;
23pub mod hybrid;
24#[cfg(feature = "replication")]
25pub mod live_replication;
26pub mod observation;
27#[cfg(feature = "root-attachment")]
28pub mod pairing;
29#[cfg(any(feature = "native", feature = "replication", feature = "iroh"))]
30pub mod publication;
31mod reopen;
32#[cfg(feature = "replication")]
33pub mod replication;
34#[cfg(all(feature = "native", feature = "iroh"))]
35pub mod replication_rpc;
36#[cfg(feature = "signing")]
37pub mod request_proof;
38#[cfg(feature = "root-attachment")]
39pub mod root_attachment;
40#[cfg(feature = "replication")]
41pub mod thread_control;
42#[cfg(feature = "replication")]
43pub mod thread_ownership;
44pub mod transport;
45
46use api::v2::client::{Client, ClientError, RpcTransport};
47pub use api::{
48 heddle::api::v1alpha2 as contract,
49 v2::{client::Rpc, rpc},
50};
51use contract::{DescribeEndpointRequest, DescribeEndpointResponse, EndpointKind, ThreadRef};
52pub use reopen::is_reopen_retryable;
53use transport::Error;
54
55pub struct Remote<T: RpcTransport<Error = Error>> {
58 pub api: Client<T>,
59 pub description: DescribeEndpointResponse,
60}
61
62impl<T: RpcTransport<Error = Error>> Remote<T> {
63 pub async fn discover(
66 transport: T,
67 endpoint_key: [u8; 32],
68 kind: EndpointKind,
69 ) -> Result<Self, ClientError<Error>> {
70 let bytes = transport
71 .unary(
72 rpc::EndpointServiceDescribeEndpoint::METHOD,
73 prost::Message::encode_to_vec(&DescribeEndpointRequest {
74 understood_packages: vec!["heddle.api.v1alpha2".into()],
75 }),
76 )
77 .await
78 .map_err(ClientError::Transport)?;
79 let description: DescribeEndpointResponse = prost::Message::decode(bytes.as_slice())?;
80 if description
81 .endpoint
82 .as_ref()
83 .is_none_or(|source| source.public_key != endpoint_key || source.kind != kind as i32)
84 || !description
85 .supported_packages
86 .iter()
87 .any(|p| p == "heddle.api.v1alpha2")
88 {
89 return Err(ClientError::Transport(Error::Protocol(
90 "endpoint identity/package mismatch",
91 )));
92 }
93 let mut api = Client::new(transport, description.implemented_methods.clone());
94 if let Some(protocol) = &description.protocol {
95 api = api.with_protocol(protocol.clone());
96 }
97 Ok(Self { api, description })
98 }
99
100 pub async fn observe<M>(
103 &self,
104 mut request: M::Request,
105 resume: Option<observation::Resume>,
106 ) -> Result<observation::Observation<T::Reader, M::Response>, observation::Error>
107 where
108 M: api::v2::client::ServerStreamingRpc,
109 M::Request: observation::ObservationRequest,
110 M::Response: observation::ObservedEvent,
111 {
112 use observation::ObservationRequest as _;
113
114 use crate::reopen::ReopenRetryable as _;
115 let budget = observation::budget(&self.description)?;
116 let options = request.options_mut();
117 options.budget = Some(budget);
118 options.after_cursor.clear();
119 let mut query = b"heddle-observation-query-v2\0".to_vec();
120 query.extend_from_slice(M::METHOD.path.as_bytes());
121 query.push(0);
122 query.extend_from_slice(&prost::Message::encode_to_vec(&request));
123 observation::validate_resume(&resume, &self.description, &query)?;
124 if let Some(resume) = &resume {
125 request.options_mut().after_cursor = resume.cursor.clone();
126 }
127 let mut attempt = 0;
128 loop {
129 let messages = match self.api.observe::<M>(&request).await {
130 Ok(messages) => messages,
131 Err(error)
132 if reopen::client_error_is_reopen_retryable(&error)
133 && attempt + 1 < reopen::ATTEMPTS =>
134 {
135 attempt += 1;
136 reopen::backoff(attempt).await;
137 continue;
138 }
139 Err(error) => return Err(error.into()),
140 };
141 let mut observation = observation::Observation::new(
142 messages,
143 &self.description,
144 budget,
145 resume.clone(),
146 query.clone(),
147 )?;
148 match observation.consume_open().await {
149 Ok(()) => return Ok(observation),
150 Err(error) if error.is_reopen_retryable() && attempt + 1 < reopen::ATTEMPTS => {
151 attempt += 1;
152 reopen::backoff(attempt).await;
153 }
154 Err(error) => {
155 observation.prime_error(error);
156 return Ok(observation);
157 }
158 }
159 }
160 }
161
162 pub async fn observe_analysis(
165 &self,
166 request: contract::ObserveAnalysisRequest,
167 resume: Option<observation::Resume>,
168 ) -> Result<observation::AnalysisObservation<T::Reader>, observation::Error> {
169 self.observe::<rpc::AnalysisServiceObserveAnalysis>(request, resume)
170 .await
171 }
172
173 pub fn thread(&self, thread: ThreadRef) -> Thread<'_, T> {
176 Thread {
177 remote: self,
178 reference: thread,
179 }
180 }
181}
182
183pub struct Thread<'a, T: RpcTransport<Error = Error>> {
184 remote: &'a Remote<T>,
185 pub reference: ThreadRef,
186}
187
188impl<T: RpcTransport<Error = Error>> Thread<'_, T> {
189 #[cfg(feature = "replication")]
192 pub async fn revise_intent(
193 &self,
194 command: &thread_control::PreparedControl,
195 ) -> Result<contract::ThreadMutationResponse, ClientError<Error>> {
196 let request = command.revise_intent().map_err(ClientError::Transport)?;
197 if request.thread.as_ref() != Some(&self.reference) {
198 return Err(ClientError::Transport(Error::Protocol(
199 "prepared command belongs to another Thread",
200 )));
201 }
202 self.remote
203 .api
204 .call::<rpc::ThreadServiceReviseIntent>(&request)
205 .await
206 }
207
208 pub async fn observe(
209 &self,
210 sections: &[contract::ThreadSection],
211 mode: contract::ObservationMode,
212 resume: Option<observation::Resume>,
213 ) -> Result<observation::ThreadObservation<T::Reader>, observation::Error> {
214 self.remote
215 .observe::<rpc::ThreadServiceObserveThread>(
216 contract::ObserveThreadRequest {
217 thread: Some(self.reference.clone()),
218 sections: sections.iter().map(|section| *section as i32).collect(),
219 observe: Some(contract::ObserveOptions {
220 mode: mode as i32,
221 ..Default::default()
222 }),
223 ..Default::default()
224 },
225 resume,
226 )
227 .await
228 }
229}