Skip to main content

heddle_thread_api/publication/
source.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Build and retain an exact selected-source transfer for explicit publication
3//! and retries. The application chooses the scratch directory and disclosure
4//! policy; preparation never reads historical States or adjacent checkouts.
5use std::{
6    fs::{File, OpenOptions},
7    path::Path,
8};
9
10use api::v2::client::RpcTransport;
11use heddle_object_model::object::{
12    EntryRedactions, ObjectSource, State, StateId, source_target::capture::ReferenceProof,
13};
14use heddle_pack::store::pack::{
15    StreamingPackBuilder, build_source_pack_with_references, build_visible_source_pack,
16};
17
18use super::{Error, PreparedPublication, PublicationOriginals};
19use crate::{Thread, contract::*, transport};
20
21pub struct SourceBudget {
22    pub max_decoded_bytes: u64,
23}
24
25/// The caller retains the same options for retry. The source endpoint is the
26/// local Iroh identity; credential signing may use a distinct owner device key.
27pub struct PublicationOptions {
28    pub client_operation_id: String,
29    pub source: EndpointRef,
30    /// Empty omits the compare-and-set condition; otherwise exactly 32 bytes.
31    /// Explicit publication does not enable ongoing synchronization.
32    pub sharing_policy_version: Vec<u8>,
33    pub checkpoint: Option<TransferCheckpoint>,
34}
35
36pub struct SourcePack {
37    directory: heddle_pack::store::pack::ScratchDir,
38    revision: StateId,
39    artifacts: [PackExtent; 2],
40}
41
42/// Read-side disclosure artifacts. A partial projection cannot be passed to
43/// `publish_source`, which requires a complete [`SourcePack`].
44pub struct VisibleSourcePack {
45    source: SourcePack,
46    complete: bool,
47}
48
49impl VisibleSourcePack {
50    /// Prepare an admitted source read, omitting hidden source entries and
51    /// reference descriptors when an entry restriction applies. Full reads
52    /// retain their independently verified reference closures.
53    pub fn prepare(
54        source: &impl ObjectSource,
55        selected: &State,
56        references: &[ReferenceProof],
57        redactions: &EntryRedactions,
58        scratch_root: &Path,
59        budget: SourceBudget,
60    ) -> Result<Self, Error> {
61        let (source, complete) = SourcePack::prepare_disclosure(
62            source,
63            selected,
64            references,
65            Some(redactions),
66            scratch_root,
67            budget,
68        )?;
69        Ok(Self { source, complete })
70    }
71
72    /// Whether the pack proves full source and selected descriptor availability.
73    pub fn is_complete(&self) -> bool {
74        self.complete
75    }
76
77    /// Whole independently hashed pack/index artifacts in transmission order.
78    pub fn artifacts(&self) -> &[PackExtent; 2] {
79        self.source.artifacts()
80    }
81
82    /// Stream the retained read artifacts while this value owns their scratch.
83    pub async fn open_artifacts(&self) -> Result<[tokio::fs::File; 2], Error> {
84        self.source.open_artifacts().await
85    }
86}
87
88impl SourcePack {
89    /// Performs disk I/O and compression. Async applications run preparation on
90    /// their blocking-work executor. The output owns and removes its scratch.
91    pub fn prepare(
92        source: &impl ObjectSource,
93        selected: &State,
94        scratch_root: &Path,
95        budget: SourceBudget,
96    ) -> Result<Self, Error> {
97        Self::prepare_with_references(source, selected, &[], scratch_root, budget)
98    }
99
100    /// Include only descriptor closures selected by independently verified source
101    /// operation proofs. A reference never authorizes another source revision.
102    pub fn prepare_with_references(
103        source: &impl ObjectSource,
104        selected: &State,
105        references: &[ReferenceProof],
106        scratch_root: &Path,
107        budget: SourceBudget,
108    ) -> Result<Self, Error> {
109        Self::prepare_disclosure(source, selected, references, None, scratch_root, budget)
110            .map(|(source, _)| source)
111    }
112
113    fn prepare_disclosure(
114        source: &impl ObjectSource,
115        selected: &State,
116        references: &[ReferenceProof],
117        redactions: Option<&EntryRedactions>,
118        scratch_root: &Path,
119        budget: SourceBudget,
120    ) -> Result<(Self, bool), Error> {
121        let directory = heddle_pack::store::pack::ScratchDir::new(scratch_root, "thread-source-")?;
122        let pack_path = directory.path().join("source.pack");
123        let index_path = directory.path().join("source.idx");
124        let pack = OpenOptions::new()
125            .read(true)
126            .write(true)
127            .create_new(true)
128            .open(&pack_path)?;
129        let builder = StreamingPackBuilder::new(
130            pack,
131            index_path.clone(),
132            Default::default(),
133            directory.path().join("buckets"),
134        )
135        .map_err(store_error)?;
136        let (pack, _, complete) = match redactions {
137            Some(redactions) => build_visible_source_pack(
138                builder,
139                source,
140                selected,
141                references,
142                redactions,
143                budget.max_decoded_bytes,
144            ),
145            None => build_source_pack_with_references(
146                builder,
147                source,
148                selected,
149                references,
150                budget.max_decoded_bytes,
151            )
152            .map(|(output, stats)| (output, stats, true)),
153        }
154        .map_err(store_error)?;
155        drop(pack);
156        let artifacts = [
157            artifact(&pack_path, pack_extent::Kind::NativePack)?,
158            artifact(&index_path, pack_extent::Kind::NativeIndex)?,
159        ];
160        Ok((
161            Self {
162                directory,
163                revision: selected.id(),
164                artifacts,
165            },
166            complete,
167        ))
168    }
169
170    /// Whole, independently hashed artifacts in transmission order. Source
171    /// preparation bounds the decoded closure before producing this inventory.
172    pub fn artifacts(&self) -> &[PackExtent; 2] {
173        &self.artifacts
174    }
175
176    /// Open the exact prepared artifacts without buffering their bodies. Keep
177    /// this SourcePack alive until the readers finish; dropping it removes its
178    /// temporary files, including after a cancelled transfer.
179    pub async fn open_artifacts(&self) -> Result<[tokio::fs::File; 2], Error> {
180        Ok([
181            tokio::fs::File::open(self.directory.path().join("source.pack")).await?,
182            tokio::fs::File::open(self.directory.path().join("source.idx")).await?,
183        ])
184    }
185
186    pub fn revision(&self) -> StateId {
187        self.revision
188    }
189
190    /// The same inventory digest verified by the publication receipt.
191    pub fn inventory_digest(&self) -> Result<[u8; 32], Error> {
192        super::inventory_digest(&self.artifacts)
193    }
194}
195
196impl<T: RpcTransport<Error = transport::Error>> Thread<'_, T> {
197    /// One exchange, directly to this Thread's endpoint. The source capture must
198    /// carry its complete original proofs. Preparation and binding are local;
199    /// no head lookup, proxy call or implicit capture happens here.
200    pub async fn publish_source(
201        &self,
202        source: &SourcePack,
203        originals: &PublicationOriginals,
204        options: PublicationOptions,
205    ) -> Result<PublicationReceipt, Error> {
206        let opening = self.publication_opening(source, options, originals)?;
207        let [pack, index] = source.open_artifacts().await?;
208        self.remote
209            .publish_content(&opening, originals, [pack, index])
210            .await
211    }
212    /// Local preparation over exact source/originals; no network authorization or
213    /// implicit acceptance. Call sign_acceptance explicitly, then send_prepared.
214    pub fn prepare_publication(
215        &self,
216        source: &SourcePack,
217        originals: PublicationOriginals,
218        options: PublicationOptions,
219        spool_genesis: heddle_object_model::object::ContentHash,
220    ) -> Result<PreparedPublication, Error> {
221        Ok(PreparedPublication::new(
222            self.publication_opening(source, options, &originals)?,
223            originals,
224            spool_genesis,
225        )?)
226    }
227
228    pub async fn send_prepared(
229        &self,
230        source: &SourcePack,
231        prepared: &PreparedPublication,
232    ) -> Result<PublicationReceipt, Error> {
233        let Some(publish_content_client_frame::Body::Open(open)) = &prepared.opening().body else {
234            return Err(Error::Invalid("prepared Open required"));
235        };
236        if open.thread.as_ref() != Some(&self.reference)
237            || open.packs.as_slice() != source.artifacts()
238            || prepared.plan().intent().revision != source.revision()
239        {
240            return Err(Error::Invalid(
241                "prepared publication differs from selected source",
242            ));
243        }
244        let artifacts = source.open_artifacts().await?;
245        self.remote
246            .publish_content(prepared.opening(), prepared.originals(), artifacts)
247            .await
248    }
249
250    fn publication_opening(
251        &self,
252        source: &SourcePack,
253        options: PublicationOptions,
254        originals: &PublicationOriginals,
255    ) -> Result<PublishContentClientFrame, Error> {
256        if self
257            .reference
258            .spool
259            .as_ref()
260            .is_none_or(|spool| spool.id.is_empty())
261            || self
262                .reference
263                .id
264                .as_ref()
265                .is_none_or(|id| id.value.len() != 32)
266            || options.source.public_key.len() != 32
267            || (!options.sharing_policy_version.is_empty()
268                && options.sharing_policy_version.len() != 32)
269        {
270            return Err(Error::Invalid(
271                "Thread, source endpoint and optional 32-byte policy version required",
272            ));
273        }
274        let destination = self
275            .remote
276            .description
277            .endpoint
278            .clone()
279            .ok_or(Error::Invalid("remote endpoint identity missing"))?;
280        let mut import_authority = None;
281        for bundle in originals
282            .operations
283            .iter()
284            .filter_map(|batch| batch.import_authority.as_ref())
285        {
286            if import_authority.is_some_and(|previous| previous != bundle) {
287                return Err(Error::Invalid(
288                    "publication originals have different import histories",
289                ));
290            }
291            crate::hybrid::bundle(Some(bundle)).map_err(Error::Invalid)?;
292            import_authority = Some(bundle);
293        }
294        let mut native_authority = None;
295        for bundle in originals
296            .operations
297            .iter()
298            .filter_map(|batch| batch.native_authority.as_ref())
299        {
300            if native_authority.is_some_and(|previous| previous != bundle) {
301                return Err(Error::Invalid(
302                    "publication originals have different native histories",
303                ));
304            }
305            native_authority = Some(bundle);
306        }
307        crate::hybrid::bundles(import_authority, native_authority).map_err(Error::Invalid)?;
308        let protocol = if import_authority.is_some()
309            || native_authority.is_some()
310            || originals
311                .geneses
312                .iter()
313                .any(|g| g.native_genesis_authority.is_some())
314        {
315            api::import_authority::require_hybrid_peer(self.remote.description.protocol.as_ref())
316                .map_err(|_| {
317                Error::Invalid(
318                    "HYBRID originals require a capable publication peer (api#307 cutover)",
319                )
320            })?;
321            Some(crate::hybrid::protocol())
322        } else {
323            crate::hybrid::sync_protocol()
324        };
325        Ok(PublishContentClientFrame {
326            client_operation_id: options.client_operation_id,
327            body: Some(publish_content_client_frame::Body::Open(
328                PublishContentOpen {
329                    native_authority: native_authority.cloned(),
330                    thread: Some(self.reference.clone()),
331                    revision: Some(RevisionRef {
332                        spool: self.reference.spool.clone(),
333                        revision: Some(revision_ref::Revision::State(
334                            api::heddle::api::common::StateId {
335                                value: source.revision.as_bytes().to_vec(),
336                            },
337                        )),
338                    }),
339                    sharing_policy_version: options.sharing_policy_version,
340                    packs: source.artifacts.to_vec(),
341                    checkpoint: options.checkpoint,
342                    source: Some(options.source),
343                    destination: Some(destination),
344                    semantic_indexes: Vec::new(),
345                    protocol,
346                    import_authority: import_authority.cloned(),
347                },
348            )),
349        })
350    }
351}
352
353fn artifact(path: &Path, kind: pack_extent::Kind) -> Result<PackExtent, Error> {
354    let mut file = File::open(path)?;
355    let length = file.metadata()?.len();
356    let mut hash = blake3::Hasher::new();
357    hash.update_reader(&mut file)?;
358    let address = ObjectAddress {
359        algorithm: "blake3".into(),
360        digest: hash.finalize().as_bytes().to_vec(),
361    };
362    Ok(PackExtent {
363        pack: Some(address.clone()),
364        kind: kind as i32,
365        offset: 0,
366        length,
367        extent_digest: Some(address),
368    })
369}
370fn store_error(error: impl std::fmt::Display) -> Error {
371    transport::Error::Io(error.to_string()).into()
372}