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 api = Client::new(transport, description.implemented_methods.clone());
94        Ok(Self { api, description })
95    }
96
97    /// Observe any contract view as atomic bounded batches. The bookmark binds
98    /// the exact method and projection as well as the authenticated endpoint.
99    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    /// Source-backed analysis uses the same committed view protocol as identity,
159    /// collaboration, checkouts and Thread observations.
160    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    /// Binding a Thread is local and costs no RPC. Persist this stable reference
170    /// when discovering/creating the Thread; its display name is never its key.
171    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    /// Send a locally prepared original signed intent. Preparation consumes the
186    /// existing overview's field frontier and performs no extra RPC.
187    #[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}