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 api = Client::new(transport, description.implemented_methods.clone());
94 Ok(Self { api, description })
95 }
96
97 pub async fn observe<M>(
100 &self,
101 mut request: M::Request,
102 resume: Option<observation::Resume>,
103 ) -> Result<observation::Observation<T::Reader, M::Response>, observation::Error>
104 where
105 M: api::v2::client::ServerStreamingRpc,
106 M::Request: observation::ObservationRequest,
107 M::Response: observation::ObservedEvent,
108 {
109 use crate::reopen::ReopenRetryable as _;
110 use observation::ObservationRequest as _;
111 let budget = observation::budget(&self.description)?;
112 let options = request.options_mut();
113 options.budget = Some(budget);
114 options.after_cursor.clear();
115 let mut query = b"heddle-observation-query-v2\0".to_vec();
116 query.extend_from_slice(M::METHOD.path.as_bytes());
117 query.push(0);
118 query.extend_from_slice(&prost::Message::encode_to_vec(&request));
119 observation::validate_resume(&resume, &self.description, &query)?;
120 if let Some(resume) = &resume {
121 request.options_mut().after_cursor = resume.cursor.clone();
122 }
123 let mut attempt = 0;
124 loop {
125 let messages = match self.api.observe::<M>(&request).await {
126 Ok(messages) => messages,
127 Err(error)
128 if reopen::client_error_is_reopen_retryable(&error)
129 && attempt + 1 < reopen::ATTEMPTS =>
130 {
131 attempt += 1;
132 reopen::backoff(attempt).await;
133 continue;
134 }
135 Err(error) => return Err(error.into()),
136 };
137 let mut observation = observation::Observation::new(
138 messages,
139 &self.description,
140 budget,
141 resume.clone(),
142 query.clone(),
143 )?;
144 match observation.consume_open().await {
145 Ok(()) => return Ok(observation),
146 Err(error) if error.is_reopen_retryable() && attempt + 1 < reopen::ATTEMPTS => {
147 attempt += 1;
148 reopen::backoff(attempt).await;
149 }
150 Err(error) => {
151 observation.prime_error(error);
152 return Ok(observation);
153 }
154 }
155 }
156 }
157
158 pub async fn observe_analysis(
161 &self,
162 request: contract::ObserveAnalysisRequest,
163 resume: Option<observation::Resume>,
164 ) -> Result<observation::AnalysisObservation<T::Reader>, observation::Error> {
165 self.observe::<rpc::AnalysisServiceObserveAnalysis>(request, resume)
166 .await
167 }
168
169 pub fn thread(&self, thread: ThreadRef) -> Thread<'_, T> {
172 Thread {
173 remote: self,
174 reference: thread,
175 }
176 }
177}
178
179pub struct Thread<'a, T: RpcTransport<Error = Error>> {
180 remote: &'a Remote<T>,
181 pub reference: ThreadRef,
182}
183
184impl<T: RpcTransport<Error = Error>> Thread<'_, T> {
185 #[cfg(feature = "replication")]
188 pub async fn revise_intent(
189 &self,
190 command: &thread_control::PreparedControl,
191 ) -> Result<contract::ThreadMutationResponse, ClientError<Error>> {
192 let request = command.revise_intent().map_err(ClientError::Transport)?;
193 if request.thread.as_ref() != Some(&self.reference) {
194 return Err(ClientError::Transport(Error::Protocol(
195 "prepared command belongs to another Thread",
196 )));
197 }
198 self.remote
199 .api
200 .call::<rpc::ThreadServiceReviseIntent>(&request)
201 .await
202 }
203
204 pub async fn observe(
205 &self,
206 sections: &[contract::ThreadSection],
207 mode: contract::ObservationMode,
208 resume: Option<observation::Resume>,
209 ) -> Result<observation::ThreadObservation<T::Reader>, observation::Error> {
210 self.remote
211 .observe::<rpc::ThreadServiceObserveThread>(
212 contract::ObserveThreadRequest {
213 thread: Some(self.reference.clone()),
214 sections: sections.iter().map(|section| *section as i32).collect(),
215 observe: Some(contract::ObserveOptions {
216 mode: mode as i32,
217 ..Default::default()
218 }),
219 ..Default::default()
220 },
221 resume,
222 )
223 .await
224 }
225}