1use 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
25pub struct PublicationOptions {
28 pub client_operation_id: String,
29 pub source: EndpointRef,
30 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
42pub struct VisibleSourcePack {
45 source: SourcePack,
46 complete: bool,
47}
48
49impl VisibleSourcePack {
50 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 pub fn is_complete(&self) -> bool {
74 self.complete
75 }
76
77 pub fn artifacts(&self) -> &[PackExtent; 2] {
79 self.source.artifacts()
80 }
81
82 pub async fn open_artifacts(&self) -> Result<[tokio::fs::File; 2], Error> {
84 self.source.open_artifacts().await
85 }
86}
87
88impl SourcePack {
89 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 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 pub fn artifacts(&self) -> &[PackExtent; 2] {
173 &self.artifacts
174 }
175
176 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 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 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 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}