Skip to main content

heddle_thread_api/
lib.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Native v2 client design. Endpoint ownership, credentials and connection
3//! discovery belong to the application; this crate never obtains a Weft token.
4#[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
55/// One authenticated source. A combined Thread retains one source per endpoint;
56/// a hosted action is sent directly to the Weft source, never via a device.
57pub 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    /// The caller supplies the Iroh-authenticated endpoint key and intended kind.
64    /// Describe is the only bootstrap exception to implemented-method discovery.
65    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    /// Observe any contract view as atomic bounded batches. The bookmark binds
101    /// the exact method and projection as well as the authenticated endpoint.
102    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    /// Source-backed analysis uses the same committed view protocol as identity,
163    /// collaboration, checkouts and Thread observations.
164    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    /// Binding a Thread is local and costs no RPC. Persist this stable reference
174    /// when discovering/creating the Thread; its display name is never its key.
175    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    /// Send a locally prepared original signed intent. Preparation consumes the
190    /// existing overview's field frontier and performs no extra RPC.
191    #[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}