Skip to main content

heddle_thread_api/
publication.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Flow-controlled source publication. Sending and receiving run together so
3//! server checkpoints cannot block an upload on a full response stream.
4use api::v2::client::{ClientError, RpcTransport};
5use prost::Message;
6use tokio::io::{AsyncRead, AsyncReadExt};
7
8use crate::{Remote, contract::*, rpc, transport};
9
10mod batching;
11
12#[cfg(feature = "source-transfer")]
13mod acceptance;
14#[cfg(feature = "source-transfer")]
15pub use acceptance::{
16    PreparedPublication, ProposedAcceptance, PublicationAcceptancePlan, proposed_publication,
17    publication_intent,
18};
19#[cfg(feature = "source-transfer")]
20mod source;
21#[cfg(feature = "source-transfer")]
22mod staging;
23#[cfg(feature = "source-transfer")]
24pub use source::{PublicationOptions, SourceBudget, SourcePack, VisibleSourcePack};
25#[cfg(feature = "source-transfer")]
26pub use staging::{
27    ProposedSourceArtifacts, validate_proposed_source_artifacts, validate_source_artifacts,
28    validate_source_artifacts_with_import_carriers,
29};
30
31/// Exact original creator wrappers and signed source operations, including
32/// every foreign integration dependency. The receiver verifies their authority
33/// and complete causal/source closure before any replica admission.
34#[derive(Clone)]
35pub struct PublicationOriginals {
36    pub geneses: Vec<ThreadGenesisRecord>,
37    pub operations: Vec<ReplicationOperations>,
38}
39impl PublicationOriginals {
40    fn validate_bounds(&self) -> Result<(), Error> {
41        if self.geneses.is_empty()
42            || self.geneses.len() > 128
43            || self.operations.is_empty()
44            || self.operations.len() > 10_000
45            || self
46                .operations
47                .iter()
48                .map(|batch| batch.operations.len())
49                .sum::<usize>()
50                > 10_000
51            || self.operations.iter().any(|batch| {
52                batch.operations.is_empty()
53                    || batch.authority_admissions.len() > batch.operations.len()
54            })
55        {
56            return Err(Error::Invalid(
57                "bounded original genesis and source operations required",
58            ));
59        }
60        let mut acceptances = std::collections::BTreeSet::new();
61        for record in self
62            .geneses
63            .iter()
64            .flat_map(|value| &value.boundary_acceptances)
65            .chain(
66                self.operations
67                    .iter()
68                    .flat_map(|value| &value.boundary_acceptances),
69            )
70        {
71            if record.canonical_record.len() > 96 * 1024
72                || record.signatures.len() != 1
73                || record.signatures[0].signature.len() != 64
74                || record.signatures[0].public_key.len() != 32
75            {
76                return Err(Error::Invalid("boundary evidence shape exceeds bounds"));
77            }
78            acceptances.insert(record.canonical_record.as_slice());
79            if acceptances.len() > 128 {
80                return Err(Error::Invalid("boundary acceptance count exceeded"));
81            }
82        }
83        let mut bytes = 0usize;
84        for length in self
85            .geneses
86            .iter()
87            .map(Message::encoded_len)
88            .chain(self.operations.iter().map(Message::encoded_len))
89        {
90            bytes = bytes
91                .checked_add(length)
92                .ok_or(Error::Invalid("publication original size overflow"))?;
93            if length > 256 * 1024 || bytes > 16 * 1024 * 1024 {
94                return Err(Error::Invalid(
95                    "publication original metadata budget exceeded",
96                ));
97            }
98        }
99        if self.geneses.iter().any(|genesis| {
100            genesis.genesis.as_ref().is_none_or(|record| {
101                record.canonical_record.is_empty() || record.signatures.is_empty()
102            })
103        }) || self
104            .operations
105            .iter()
106            .flat_map(|batch| {
107                batch
108                    .operations
109                    .iter()
110                    .chain(&batch.authority_admissions)
111                    .chain(&batch.boundary_acceptances)
112            })
113            .any(|record| record.canonical_record.is_empty() || record.signatures.is_empty())
114        {
115            return Err(Error::Invalid(
116                "original signatures and canonical records required",
117            ));
118        }
119        Ok(())
120    }
121}
122
123#[derive(Debug, thiserror::Error)]
124pub enum Error {
125    #[error(transparent)]
126    Client(#[from] ClientError<transport::Error>),
127    #[error(transparent)]
128    Transport(#[from] transport::Error),
129    #[error("publication source: {0}")]
130    Source(#[from] std::io::Error),
131    #[error("invalid publication: {0}")]
132    Invalid(&'static str),
133    #[error(
134        "publication {limit_name} limit is {limit}, required {actual} for {operations} operations"
135    )]
136    OriginalBudgetExceeded {
137        operations: usize,
138        limit_name: &'static str,
139        limit: usize,
140        actual: usize,
141    },
142    #[error(
143        "publication operation {operation} requires {bytes} encoded bytes; batch limit {batch_limit} bytes, frame limit {frame_limit} bytes"
144    )]
145    OriginalOperationTooLarge {
146        operation: usize,
147        bytes: usize,
148        batch_limit: usize,
149        frame_limit: usize,
150    },
151}
152
153impl<T: RpcTransport<Error = transport::Error>> Remote<T> {
154    /// Publish original source proofs and exact source in one exchange. Inputs are the native pack
155    /// and its index in opening order; neither is buffered in full. The caller
156    /// retains the operation ID and opening to retry after an interrupted call.
157    /// A receipt is returned only after verifying its exact publication scope.
158    pub async fn publish_content<R: AsyncRead + Unpin + Send>(
159        &self,
160        opening: &PublishContentClientFrame,
161        originals: &PublicationOriginals,
162        mut artifacts: [R; 2],
163    ) -> Result<PublicationReceipt, Error> {
164        // Re-batching retains each unit's complete original evidence.
165        for batch in &originals.operations {
166            crate::hybrid::operations(batch).map_err(Error::Invalid)?;
167        }
168        // Preflight all originals before opening the exchange. The Ready may
169        // narrow the frame budget; re-batch again before sending any content.
170        let originals =
171            batching::bounded_originals(originals, &opening.client_operation_id, 512 * 1024)?;
172        originals.validate_bounds()?;
173        let Some(publish_content_client_frame::Body::Open(open)) = &opening.body else {
174            return Err(Error::Invalid("Open required"));
175        };
176        crate::hybrid::publish_open(open).map_err(Error::Invalid)?;
177        if open.protocol.is_some() {
178            api::import_authority::require_hybrid_peer(self.description.protocol.as_ref())
179                .map_err(|_| Error::Invalid("peer does not support HYBRID publication"))?;
180        }
181        if originals.operations.iter().any(|batch| {
182            (batch.import_authority.is_some() && batch.import_authority != open.import_authority)
183                || (batch.native_authority.is_some()
184                    && batch.native_authority != open.native_authority)
185        }) {
186            return Err(Error::Invalid(
187                "publication proof differs from negotiated opening",
188            ));
189        }
190        if open.destination != self.description.endpoint
191            || open.thread.is_none()
192            || open.revision.is_none()
193            || (!open.sharing_policy_version.is_empty() && open.sharing_policy_version.len() != 32)
194            || open.packs.len() != 2
195            || open.packs[0].kind != pack_extent::Kind::NativePack as i32
196            || open.packs[1].kind != pack_extent::Kind::NativeIndex as i32
197        {
198            return Err(Error::Invalid(
199                "endpoint, capture and ordered native artifacts required",
200            ));
201        }
202        let inventory = inventory_digest(&open.packs)?;
203        let mut logical = opening.clone();
204        if let Some(publish_content_client_frame::Body::Open(open)) = logical.body.as_mut() {
205            open.checkpoint = None;
206        }
207        let digest = typed_digest("thread-source-transfer-v1", &logical.encode_to_vec());
208        let expected = TransferCheckpoint {
209            transfer_id: digest[..16].to_vec(),
210            plan_digest: digest.to_vec(),
211            resume_token: Vec::new(),
212            committed_bytes: 0,
213        };
214        let (mut sender, mut responses) = self
215            .api
216            .exchange::<rpc::SyncServicePublishContent>(opening)
217            .await?;
218        let first = responses
219            .next()
220            .await?
221            .ok_or(Error::Invalid("publication ended before admission"))?;
222        let ready = match first.body {
223            Some(publish_content_server_frame::Body::Receipt(receipt)) => {
224                return validate_receipt(receipt, opening, &inventory);
225            }
226            Some(publish_content_server_frame::Body::Ready(ready)) => ready,
227            _ => return Err(Error::Invalid("Ready or replay receipt required")),
228        };
229        crate::hybrid::transfer_ready(&ready).map_err(Error::Invalid)?;
230        crate::hybrid::negotiated(open.protocol.as_ref(), ready.protocol.as_ref())
231            .map_err(Error::Invalid)?;
232        if ready.endpoint != open.destination
233            || ready.thread != open.thread
234            || ready.current != open.revision
235            || ready.checkpoint.as_ref() != Some(&expected)
236        {
237            return Err(Error::Invalid("admission differs from publication plan"));
238        }
239        let budget = ready
240            .budget
241            .ok_or(Error::Invalid("publication frame budget required"))?;
242        let frame_limit = budget.max_frame_bytes as usize;
243        if !(1024..=512 * 1024).contains(&frame_limit) {
244            return Err(Error::Invalid("unsupported publication frame budget"));
245        }
246        let originals =
247            batching::bounded_originals(&originals, &opening.client_operation_id, frame_limit)?;
248        originals.validate_bounds()?;
249        let upload = async {
250            for body in originals
251                .geneses
252                .iter()
253                .cloned()
254                .map(publish_content_client_frame::Body::ThreadGenesis)
255                .chain(
256                    originals
257                        .operations
258                        .iter()
259                        .cloned()
260                        .map(publish_content_client_frame::Body::Operations),
261                )
262            {
263                let frame = PublishContentClientFrame {
264                    client_operation_id: opening.client_operation_id.clone(),
265                    body: Some(body),
266                };
267                if frame.encoded_len() > frame_limit {
268                    return Err(Error::Invalid(
269                        "original exceeds negotiated publication frame budget",
270                    ));
271                }
272                sender.send(&frame).await?;
273            }
274            for (artifact, planned) in artifacts.iter_mut().zip(&open.packs) {
275                let mut offset = 0;
276                let mut digest = blake3::Hasher::new();
277                let mut buffer = vec![0; frame_limit / 2];
278                while offset < planned.length {
279                    let length = (planned.length - offset).min(buffer.len() as u64) as usize;
280                    let read = artifact.read(&mut buffer[..length]).await?;
281                    if read == 0 {
282                        return Err(Error::Invalid("artifact ended before declared length"));
283                    }
284                    let data = &buffer[..read];
285                    digest.update(data);
286                    let chunk = PublishContentClientFrame {
287                        client_operation_id: opening.client_operation_id.clone(),
288                        body: Some(publish_content_client_frame::Body::Pack(PackChunk {
289                            extent: Some(PackExtent {
290                                pack: planned.pack.clone(),
291                                kind: planned.kind,
292                                offset,
293                                length: read as u64,
294                                extent_digest: Some(ObjectAddress {
295                                    algorithm: "blake3".into(),
296                                    digest: blake3::hash(data).as_bytes().to_vec(),
297                                }),
298                            }),
299                            data: data.to_vec(),
300                        })),
301                    };
302                    if chunk.encoded_len() > frame_limit {
303                        return Err(Error::Invalid("chunk exceeds negotiated frame budget"));
304                    }
305                    sender.send(&chunk).await?;
306                    offset += read as u64;
307                }
308                if artifact.read(&mut buffer[..1]).await? != 0
309                    || planned
310                        .pack
311                        .as_ref()
312                        .is_none_or(|address| address.digest != digest.finalize().as_bytes())
313                {
314                    return Err(Error::Invalid(
315                        "artifact differs from declared length or digest",
316                    ));
317                }
318            }
319            sender
320                .send(&PublishContentClientFrame {
321                    client_operation_id: opening.client_operation_id.clone(),
322                    body: Some(publish_content_client_frame::Body::Finish(
323                        PublishContentFinish {
324                            checkpoint: Some(expected.clone()),
325                        },
326                    )),
327                })
328                .await?;
329            sender.finish().await?;
330            Ok::<_, Error>(())
331        };
332        let receive = async {
333            while let Some(frame) = responses.next().await? {
334                match frame.body {
335                    Some(publish_content_server_frame::Body::Checkpoint(checkpoint))
336                        if checkpoint == expected => {}
337                    Some(publish_content_server_frame::Body::Receipt(receipt)) => {
338                        return validate_receipt(receipt, opening, &inventory);
339                    }
340                    _ => return Err(Error::Invalid("unexpected publication response")),
341                }
342            }
343            Err(Error::Invalid(
344                "publication ended without a durable receipt",
345            ))
346        };
347        let (_, receipt) = tokio::try_join!(upload, receive)?;
348        Ok(receipt)
349    }
350}
351
352fn validate_receipt(
353    receipt: PublicationReceipt,
354    opening: &PublishContentClientFrame,
355    inventory: &[u8; 32],
356) -> Result<PublicationReceipt, Error> {
357    let Some(publish_content_client_frame::Body::Open(open)) = &opening.body else {
358        return Err(Error::Invalid("Open required"));
359    };
360    crate::hybrid::publication_receipt(&receipt).map_err(Error::Invalid)?;
361    if receipt.import_authority.is_some() || receipt.native_authority.is_some() {
362        api::import_authority::require_hybrid_peer(open.protocol.as_ref())
363            .map_err(|_| Error::Invalid("HYBRID receipt requires negotiated publication"))?;
364    }
365    if receipt.client_operation_id != opening.client_operation_id
366        || receipt.destination != open.destination
367        || receipt.thread != open.thread
368        || receipt.revision != open.revision
369        || receipt.sharing_policy_version.len() != 32
370        || (!open.sharing_policy_version.is_empty()
371            && receipt.sharing_policy_version != open.sharing_policy_version)
372    {
373        return Err(Error::Invalid("receipt differs from requested publication"));
374    }
375    match &receipt.outcome {
376        Some(publication_receipt::Outcome::Accepted(_))
377            if receipt.accepted_inventory.as_ref().is_some_and(|address| {
378                address.algorithm == "blake3" && address.digest == inventory
379            }) =>
380        {
381            Ok(receipt)
382        }
383        Some(publication_receipt::Outcome::Rejected(error)) => {
384            Err(transport::Error::Remote(error.clone().into()).into())
385        }
386        _ => Err(Error::Invalid(
387            "receipt did not accept the complete inventory",
388        )),
389    }
390}
391
392/// Exact ordered complete artifact inventory shared by preparation, intent,
393/// upload, and receipt verification. No pack body is buffered here.
394pub fn inventory_digest(packs: &[PackExtent]) -> Result<[u8; 32], Error> {
395    if packs.len() != 2
396        || packs[0].kind != pack_extent::Kind::NativePack as i32
397        || packs[1].kind != pack_extent::Kind::NativeIndex as i32
398    {
399        return Err(Error::Invalid("ordered native pack and index required"));
400    }
401    let mut inventory = Vec::new();
402    let mut total = 0_u64;
403    for extent in packs {
404        let address = extent
405            .pack
406            .as_ref()
407            .ok_or(Error::Invalid("artifact address required"))?;
408        total = total
409            .checked_add(extent.length)
410            .ok_or(Error::Invalid("artifact length overflow"))?;
411        if address.algorithm != "blake3"
412            || address.digest.len() != 32
413            || extent.length == 0
414            || extent.offset != 0
415            || extent.extent_digest.as_ref() != Some(address)
416            || total > 256 * 1024 * 1024
417        {
418            return Err(Error::Invalid("complete bounded BLAKE3 artifacts required"));
419        }
420        extent
421            .encode_length_delimited(&mut inventory)
422            .map_err(|_| Error::Invalid("inventory encoding failed"))?;
423    }
424    Ok(typed_digest("thread-source-inventory-v1", &inventory))
425}
426
427fn typed_digest(kind: &str, bytes: &[u8]) -> [u8; 32] {
428    let mut hasher = blake3::Hasher::new();
429    hasher.update(kind.as_bytes());
430    hasher.update(&(bytes.len() as u64).to_le_bytes());
431    hasher.update(&[0]);
432    hasher.update(bytes);
433    *hasher.finalize().as_bytes()
434}
435
436#[cfg(test)]
437mod tests {
438    use api::v2::{
439        MethodDescriptor,
440        client::{MessageReader, MessageWriter},
441    };
442    use tokio::sync::mpsc;
443
444    use super::*;
445
446    pub(super) struct Reader(mpsc::Receiver<Vec<u8>>);
447    pub(super) struct Writer(Option<mpsc::Sender<Vec<u8>>>);
448    impl MessageReader for Reader {
449        type Error = transport::Error;
450        async fn next(&mut self) -> Result<Option<Vec<u8>>, Self::Error> {
451            Ok(self.0.recv().await)
452        }
453        fn cancel(&mut self) {
454            self.0.close();
455        }
456    }
457    impl MessageWriter for Writer {
458        type Error = transport::Error;
459        async fn send(&mut self, bytes: Vec<u8>) -> Result<(), Self::Error> {
460            self.0
461                .as_ref()
462                .ok_or(transport::Error::Protocol("closed"))?
463                .send(bytes)
464                .await
465                .map_err(|_| transport::Error::Protocol("peer closed"))
466        }
467        async fn finish(&mut self) -> Result<(), Self::Error> {
468            self.0.take();
469            Ok(())
470        }
471        fn abort(&mut self) {
472            self.0.take();
473        }
474    }
475    pub(super) struct Peer {
476        wrong_receipt: bool,
477        original_count: usize,
478    }
479    impl RpcTransport for Peer {
480        type Error = transport::Error;
481        type Reader = Reader;
482        type Writer = Writer;
483        async fn unary(
484            &self,
485            _: &'static MethodDescriptor,
486            _: Vec<u8>,
487        ) -> Result<Vec<u8>, Self::Error> {
488            unreachable!("test exercises only publication");
489        }
490        async fn observe(
491            &self,
492            _: &'static MethodDescriptor,
493            _: Vec<u8>,
494        ) -> Result<Reader, Self::Error> {
495            unreachable!("test exercises only publication");
496        }
497        async fn exchange(
498            &self,
499            _: &'static MethodDescriptor,
500            bytes: Vec<u8>,
501        ) -> Result<(Writer, Reader), Self::Error> {
502            let opening = PublishContentClientFrame::decode(bytes.as_slice())?;
503            let Some(publish_content_client_frame::Body::Open(open)) = opening.body.clone() else {
504                panic!("Open");
505            };
506            let mut logical = opening.clone();
507            if let Some(publish_content_client_frame::Body::Open(open)) = logical.body.as_mut() {
508                open.checkpoint = None;
509            }
510            let digest = typed_digest("thread-source-transfer-v1", &logical.encode_to_vec());
511            let checkpoint = TransferCheckpoint {
512                transfer_id: digest[..16].to_vec(),
513                plan_digest: digest.to_vec(),
514                committed_bytes: 0,
515                resume_token: vec![],
516            };
517            let (tx, mut incoming) = mpsc::channel::<Vec<u8>>(1);
518            let (outgoing, rx) = mpsc::channel(1);
519            let wrong_receipt = self.wrong_receipt;
520            let original_count = self.original_count;
521            tokio::spawn(async move {
522                let send = |body| {
523                    let outgoing = &outgoing;
524                    async move {
525                        outgoing
526                            .send(PublishContentServerFrame { body: Some(body) }.encode_to_vec())
527                            .await
528                    }
529                };
530                send(publish_content_server_frame::Body::Ready(TransferReady {
531                    endpoint: open.destination.clone(),
532                    thread: open.thread.clone(),
533                    current: open.revision.clone(),
534                    checkpoint: Some(checkpoint.clone()),
535                    protocol: open.protocol.clone(),
536                    import_authority: open.import_authority.clone(),
537                    budget: Some(ReadBudget {
538                        max_frame_bytes: if open.import_authority.is_some() {
539                            256 * 1024
540                        } else {
541                            2048
542                        },
543                        ..Default::default()
544                    }),
545                    ..Default::default()
546                }))
547                .await
548                .expect("Ready");
549                let mut lengths = [0_u64; 2];
550                let mut original_counts = [0usize; 2];
551                while let Some(bytes) = incoming.recv().await {
552                    assert!(
553                        bytes.len()
554                            <= if open.import_authority.is_some() {
555                                256 * 1024
556                            } else {
557                                2048
558                            },
559                        "negotiated frame size"
560                    );
561                    let frame = PublishContentClientFrame::decode(bytes.as_slice()).expect("frame");
562                    assert_eq!(frame.client_operation_id, opening.client_operation_id);
563                    match frame.body {
564                        Some(publish_content_client_frame::Body::ThreadGenesis(genesis)) => {
565                            assert_eq!(lengths, [0, 0]);
566                            assert!(genesis.genesis.is_some());
567                            original_counts[0] += 1;
568                            send(publish_content_server_frame::Body::Checkpoint(
569                                checkpoint.clone(),
570                            ))
571                            .await
572                            .expect("original checkpoint");
573                        }
574                        Some(publish_content_client_frame::Body::Operations(batch)) => {
575                            assert_eq!(lengths, [0, 0]);
576                            assert!((1..=128).contains(&batch.operations.len()));
577                            assert!(!batch.operations[0].canonical_record.is_empty());
578                            original_counts[1] += batch.operations.len();
579                            send(publish_content_server_frame::Body::Checkpoint(
580                                checkpoint.clone(),
581                            ))
582                            .await
583                            .expect("original checkpoint");
584                        }
585                        Some(publish_content_client_frame::Body::Pack(chunk)) => {
586                            assert_eq!(
587                                original_counts,
588                                [1, original_count],
589                                "original proofs precede source artifacts"
590                            );
591                            let extent = chunk.extent.expect("extent");
592                            let index = if extent.kind == pack_extent::Kind::NativePack as i32 {
593                                0
594                            } else {
595                                1
596                            };
597                            assert_eq!(extent.offset, lengths[index]);
598                            assert_eq!(extent.length, chunk.data.len() as u64);
599                            assert_eq!(
600                                extent.extent_digest.expect("chunk digest").digest,
601                                blake3::hash(&chunk.data).as_bytes()
602                            );
603                            lengths[index] += extent.length;
604                            // Capacity one in BOTH directions. An upload that
605                            // waits to read responses until all sends complete
606                            // deadlocks after a few chunks here.
607                            send(publish_content_server_frame::Body::Checkpoint(
608                                checkpoint.clone(),
609                            ))
610                            .await
611                            .expect("checkpoint");
612                        }
613                        Some(publish_content_client_frame::Body::Finish(finish)) => {
614                            assert_eq!(finish.checkpoint, Some(checkpoint.clone()));
615                            assert_eq!(lengths, [open.packs[0].length, open.packs[1].length]);
616                            let mut inventory = Vec::new();
617                            for extent in &open.packs {
618                                extent
619                                    .encode_length_delimited(&mut inventory)
620                                    .expect("inventory");
621                            }
622                            let mut receipt = PublicationReceipt {
623                                native_authority: None,
624                                client_operation_id: opening.client_operation_id.clone(),
625                                destination: open.destination.clone(),
626                                thread: open.thread.clone(),
627                                revision: open.revision.clone(),
628                                sharing_policy_version: if open.sharing_policy_version.is_empty() {
629                                    vec![3; 32]
630                                } else {
631                                    open.sharing_policy_version.clone()
632                                },
633                                accepted_inventory: Some(ObjectAddress {
634                                    algorithm: "blake3".into(),
635                                    digest: typed_digest("thread-source-inventory-v1", &inventory)
636                                        .to_vec(),
637                                }),
638                                outcome: Some(publication_receipt::Outcome::Accepted(
639                                    Applied::default(),
640                                )),
641                                import_authority: open.import_authority.clone(),
642                            };
643                            if wrong_receipt {
644                                receipt.thread = None;
645                            }
646                            let _ =
647                                send(publish_content_server_frame::Body::Receipt(receipt)).await;
648                            break;
649                        }
650                        _ => panic!("unexpected frame"),
651                    }
652                }
653            });
654            Ok((Writer(Some(tx)), Reader(rx)))
655        }
656    }
657
658    pub(super) fn fixture(
659        wrong_receipt: bool,
660    ) -> (
661        Remote<Peer>,
662        PublishContentClientFrame,
663        [std::io::Cursor<Vec<u8>>; 2],
664    ) {
665        let endpoint = EndpointRef {
666            public_key: vec![8; 32],
667            kind: EndpointKind::Weft as i32,
668        };
669        let remote = Remote {
670            api: api::v2::client::Client::new(
671                Peer {
672                    wrong_receipt,
673                    original_count: 1,
674                },
675                ["/heddle.api.v1alpha2.SyncService/PublishContent".into()],
676            ),
677            description: DescribeEndpointResponse {
678                endpoint: Some(endpoint.clone()),
679                ..Default::default()
680            },
681        };
682        let artifacts = [vec![1; 128 * 1024], vec![2; 16 * 1024]];
683        let packs = artifacts
684            .iter()
685            .zip([
686                pack_extent::Kind::NativePack,
687                pack_extent::Kind::NativeIndex,
688            ])
689            .map(|(data, kind)| {
690                let address = ObjectAddress {
691                    algorithm: "blake3".into(),
692                    digest: blake3::hash(data).as_bytes().to_vec(),
693                };
694                PackExtent {
695                    pack: Some(address.clone()),
696                    kind: kind as i32,
697                    offset: 0,
698                    length: data.len() as u64,
699                    extent_digest: Some(address),
700                }
701            })
702            .collect();
703        let open = PublishContentClientFrame {
704            client_operation_id: "op-test".into(),
705            body: Some(publish_content_client_frame::Body::Open(
706                PublishContentOpen {
707                    destination: Some(endpoint),
708                    thread: Some(ThreadRef::default()),
709                    revision: Some(RevisionRef::default()),
710                    packs,
711                    ..Default::default()
712                },
713            )),
714        };
715        (remote, open, artifacts.map(std::io::Cursor::new))
716    }
717    // These byte fixtures exercise framing/backpressure, not original authority
718    // admission; real Iroh tests independently verify canonical signatures.
719    fn originals() -> PublicationOriginals {
720        let record = SignedRecord {
721            format: "transport-fixture".into(),
722            canonical_record: vec![1],
723            signatures: vec![RecordSignature {
724                public_key: vec![2; 32],
725                signature: vec![3; 64],
726            }],
727        };
728        PublicationOriginals {
729            geneses: vec![ThreadGenesisRecord {
730                boundary_acceptances: Vec::new(),
731                genesis: Some(record.clone()),
732                ..Default::default()
733            }],
734            operations: vec![ReplicationOperations {
735                native_authority: None,
736                boundary_acceptances: Vec::new(),
737                operations: vec![record],
738                authority_admissions: vec![],
739                import_authority: None,
740            }],
741        }
742    }
743
744    #[test]
745    fn publication_originals_require_bounded_complete_metadata() {
746        let valid = originals();
747        assert!(valid.validate_bounds().is_ok());
748        let mut missing = valid.clone();
749        missing.operations.clear();
750        assert!(matches!(
751            missing.validate_bounds(),
752            Err(Error::Invalid(
753                "bounded original genesis and source operations required"
754            ))
755        ));
756        let mut unsigned = valid.clone();
757        unsigned.operations[0].operations[0].signatures.clear();
758        assert!(matches!(
759            unsigned.validate_bounds(),
760            Err(Error::Invalid(
761                "original signatures and canonical records required"
762            ))
763        ));
764        let mut oversized = valid.clone();
765        oversized.operations[0].operations[0].canonical_record = vec![1; 256 * 1024];
766        assert!(matches!(
767            oversized.validate_bounds(),
768            Err(Error::Invalid(
769                "publication original metadata budget exceeded"
770            ))
771        ));
772        let mut too_many = valid;
773        too_many.geneses = vec![too_many.geneses[0].clone(); 129];
774        assert!(matches!(
775            too_many.validate_bounds(),
776            Err(Error::Invalid(
777                "bounded original genesis and source operations required"
778            ))
779        ));
780    }
781
782    #[tokio::test]
783    async fn publication_drains_checkpoints_while_uploading_under_backpressure() {
784        let (remote, open, artifacts) = fixture(false);
785        let receipt = tokio::time::timeout(
786            std::time::Duration::from_secs(2),
787            remote.publish_content(&open, &originals(), artifacts),
788        )
789        .await
790        .expect("must not deadlock")
791        .expect("publication");
792        assert_eq!(receipt.client_operation_id, "op-test");
793    }
794    #[tokio::test]
795    async fn publication_rebatches_after_ready_before_upload() {
796        let (mut remote, open, artifacts) = fixture(false);
797        remote.api = api::v2::client::Client::new(
798            Peer {
799                wrong_receipt: false,
800                original_count: 140,
801            },
802            ["/heddle.api.v1alpha2.SyncService/PublishContent".into()],
803        );
804        let mut originals = originals();
805        let mut record = originals.operations[0].operations[0].clone();
806        record.canonical_record = vec![1; 300];
807        originals.operations[0].operations = vec![record; 140];
808        let receipt = tokio::time::timeout(
809            std::time::Duration::from_secs(5),
810            remote.publish_content(&open, &originals, artifacts),
811        )
812        .await
813        .expect("must not deadlock")
814        .expect("negotiated bounded publication");
815        assert_eq!(receipt.client_operation_id, "op-test");
816    }
817
818    #[cfg(feature = "native")]
819    #[tokio::test]
820    async fn hosted_source_publication_retains_complete_history_and_original_signatures() {
821        use objects::store::{FsStore, ObjectStore};
822        let scratch = tempfile::tempdir().expect("source");
823        let (staged, _, pinned) = crate::fetch::hosted::tests::source(scratch.path(), false);
824        let bundle = staged.import_authority().expect("complete history").clone();
825        let store = FsStore::new(scratch.path().join("publication-store"));
826        store.init().expect("store");
827        let paths = staged.artifact_paths();
828        store
829            .install_pack_streaming(&paths[0], &paths[1])
830            .expect("isolated original pack");
831        let pack = SourcePack::prepare(
832            &store,
833            staged.state(),
834            scratch.path(),
835            SourceBudget {
836                max_objects: 16,
837                max_decoded_bytes: 1024 * 1024,
838            },
839        )
840        .expect("selected closure");
841        let operations = crate::authority_admission::batches(
842            staged.operations().iter().cloned().map(|original| {
843                crate::replication::store::ReceivedOperation {
844                    native_authority: None,
845                    original,
846                    authority_admission: None,
847                    import_authority: Some(std::sync::Arc::new(bundle.clone())),
848                }
849            }),
850            128 * 1024,
851            128,
852        )
853        .expect("bounded carrier")
854        .collect::<Result<Vec<_>, _>>()
855        .expect("originals");
856        let originals = PublicationOriginals {
857            geneses: vec![
858                staged
859                    .ready()
860                    .thread_genesis
861                    .clone()
862                    .expect("selected genesis"),
863            ],
864            operations,
865        };
866        let selected = staged.ready().thread.clone().expect("Thread");
867        let (mut remote, _, _) = fixture(false);
868        let options = || PublicationOptions {
869            client_operation_id: "2a2a2a2a-2a2a-2a2a-2a2a-2a2a2a2a2a2a".into(),
870            source: EndpointRef {
871                public_key: vec![2; 32],
872                kind: EndpointKind::Device as i32,
873            },
874            sharing_policy_version: vec![],
875            checkpoint: None,
876        };
877        assert!(
878            remote
879                .thread(selected.clone())
880                .publish_source(&pack, &originals, options())
881                .await
882                .is_err(),
883            "old peer cannot strip history"
884        );
885        remote.description.protocol = Some(crate::hybrid::protocol());
886        let receipt = remote
887            .thread(selected)
888            .publish_source(&pack, &originals, options())
889            .await
890            .expect("capable publication");
891        assert_eq!(receipt.import_authority.as_ref(), Some(&bundle));
892        let history = crate::hybrid::authority::AcceptedHistory::from_selected_spool(
893            &bundle,
894            &pinned,
895            1350,
896            heddleco_capability_verifier::VerificationLimits::new(30 * 24 * 60 * 60)
897                .expect("limits"),
898        )
899        .expect("independent owner");
900        let prepared = remote
901            .thread(staged.ready().thread.clone().expect("Thread"))
902            .prepare_publication(
903                &pack,
904                originals.clone(),
905                PublicationOptions {
906                    client_operation_id: "2b2b2b2b-2b2b-2b2b-2b2b-2b2b2b2b2b2b".into(),
907                    source: EndpointRef {
908                        public_key: vec![2; 32],
909                        kind: EndpointKind::Device as i32,
910                    },
911                    sharing_policy_version: vec![],
912                    checkpoint: None,
913                },
914                heddle_object_model::object::ContentHash::from_bytes(*history.genesis()),
915            )
916            .expect("prepared public originals");
917        let Some(publish_content_client_frame::Body::Open(open)) = &prepared.opening().body else {
918            panic!("Open")
919        };
920        assert_eq!(open.import_authority.as_ref(), Some(&bundle));
921        let received = tempfile::tempdir().expect("received artifacts");
922        for (mut file, name) in pack
923            .open_artifacts()
924            .await
925            .expect("original artifacts")
926            .into_iter()
927            .zip(["source.pack", "source.idx"])
928        {
929            let mut output = tokio::fs::File::create(received.path().join(name))
930                .await
931                .expect("received file");
932            tokio::io::copy(&mut file, &mut output)
933                .await
934                .expect("exact uploaded bytes");
935        }
936        let authority = crate::hybrid::authority::SelectedAuthority::new(
937            history,
938            bundle.clone(),
939            |_: &ImportPublicProofBundleV1,
940             _: i64,
941             _: &repo::thread_replication::hosted_trust::TrustTransaction<'_>| Ok(()),
942        );
943        let fixture: serde_json::Value =
944            serde_json::from_str(include_str!("../tests/fixtures/hybrid-alpha33.json"))
945                .expect("tagged vectors");
946        let carriers = repo::thread_replication::delegated_import::authenticate_import_carriers(
947            &bundle,
948            &authority,
949            &api::import_authority::ImportWitnessRootPin {
950                authority: "https://weft.example.test".into(),
951                root_id: "descriptor-root-1".into(),
952                public_key: hex::decode(
953                    fixture["keys"]["root"]["public_key_hex"]
954                        .as_str()
955                        .expect("root"),
956                )
957                .expect("root key"),
958                epoch: 1,
959            },
960            1_350_000,
961            &[],
962            &[],
963            |_| Ok(()),
964        )
965        .expect("independently authenticated import carriers");
966        let validated = validate_source_artifacts_with_import_carriers(
967            received,
968            open,
969            prepared.originals().clone(),
970            carriers,
971        )
972        .expect("signed originals and actual closure");
973        assert_eq!(validated.import_authority(), Some(&bundle));
974        let received = validated
975            .into_hosted_source(staged.ready().clone())
976            .expect("retain public history through staging");
977        assert_eq!(received.import_authority(), Some(&bundle));
978        assert_eq!(received.operations(), staged.operations());
979
980        assert_eq!(
981            originals.operations[0].operations[0].canonical_record,
982            staged.operations()[0].canonical
983        );
984        assert_eq!(
985            originals.operations[0].operations[0].signatures[0].signature,
986            staged.operations()[0].signature
987        );
988        assert_eq!(
989            bundle.original_geneses.len(),
990            2,
991            "complete sibling public history survives selected source publication"
992        );
993    }
994
995    #[tokio::test]
996    async fn publication_refuses_originals_carrying_hybrid_import_authority() {
997        let (remote, open, artifacts) = fixture(false);
998        let mut originals = originals();
999        originals.operations[0].import_authority = Some(ImportPublicProofBundleV1::default());
1000        let error = remote
1001            .publish_content(&open, &originals, artifacts)
1002            .await
1003            .expect_err("HYBRID originals must not be relayed with the bundle dropped");
1004        assert!(
1005            matches!(error, Error::Invalid(message) if message.contains("api#307")),
1006            "{error}"
1007        );
1008    }
1009
1010    #[tokio::test]
1011    async fn publication_refuses_a_receipt_for_another_scope() {
1012        let (remote, open, artifacts) = fixture(true);
1013        let error = remote
1014            .publish_content(&open, &originals(), artifacts)
1015            .await
1016            .expect_err("mismatched scope");
1017        assert!(matches!(
1018            error,
1019            Error::Invalid("receipt differs from requested publication")
1020        ));
1021    }
1022    #[test]
1023    fn publication_policy_is_optional_cas_and_receipt_reports_actual_frontier() {
1024        let (_, mut opening, _) = fixture(false);
1025        let inventory = [7; 32];
1026        let Some(publish_content_client_frame::Body::Open(open)) = &opening.body else {
1027            panic!("opening")
1028        };
1029        let receipt = PublicationReceipt {
1030            native_authority: None,
1031            client_operation_id: opening.client_operation_id.clone(),
1032            destination: open.destination.clone(),
1033            thread: open.thread.clone(),
1034            revision: open.revision.clone(),
1035            sharing_policy_version: vec![3; 32],
1036            accepted_inventory: Some(ObjectAddress {
1037                algorithm: "blake3".into(),
1038                digest: inventory.to_vec(),
1039            }),
1040            outcome: Some(publication_receipt::Outcome::Accepted(Applied::default())),
1041            import_authority: None,
1042        };
1043        validate_receipt(receipt.clone(), &opening, &inventory)
1044            .expect("one-shot upload returns actual policy without CAS");
1045        let mut hybrid = receipt.clone();
1046        hybrid.import_authority = Some(ImportPublicProofBundleV1::default());
1047        assert!(
1048            matches!(
1049                validate_receipt(hybrid, &opening, &inventory),
1050                Err(Error::Invalid(message)) if message.contains("api#307")
1051            ),
1052            "a receipt carrying HYBRID import authority must be refused, never ignored"
1053        );
1054        let mut absent = receipt.clone();
1055        absent.sharing_policy_version.clear();
1056        assert!(
1057            validate_receipt(absent, &opening, &inventory).is_err(),
1058            "receipt must report actual policy frontier"
1059        );
1060        let Some(publish_content_client_frame::Body::Open(open)) = &mut opening.body else {
1061            panic!("opening")
1062        };
1063        open.sharing_policy_version = vec![4; 32];
1064        assert!(
1065            validate_receipt(receipt.clone(), &opening, &inventory).is_err(),
1066            "explicit policy CAS cannot silently accept another version"
1067        );
1068        let Some(publish_content_client_frame::Body::Open(open)) = &mut opening.body else {
1069            panic!("opening")
1070        };
1071        open.sharing_policy_version = vec![3; 32];
1072        validate_receipt(receipt, &opening, &inventory).expect("matching explicit policy version");
1073    }
1074
1075    #[cfg(feature = "replication")]
1076    #[test]
1077    fn publication_digests_match_the_native_typed_hash_format() {
1078        assert_eq!(
1079            typed_digest("thread-source-inventory-v1", b"bytes"),
1080            *heddle_object_model::object::ContentHash::compute_typed(
1081                "thread-source-inventory-v1",
1082                b"bytes"
1083            )
1084            .as_bytes()
1085        );
1086    }
1087
1088    #[cfg(feature = "source-transfer")]
1089    #[tokio::test]
1090    async fn thread_publication_prepares_only_selected_source_and_binds_its_revision() {
1091        use objects::{
1092            object::{Attribution, Blob, Principal, State, Tree, TreeEntry},
1093            store::{FsStore, ObjectStore},
1094        };
1095        let root = tempfile::tempdir().expect("local source scratch");
1096        let store = FsStore::new(root.path().join("objects"));
1097        store.init().expect("store");
1098        let blob = Blob::new(vec![1; 128 * 1024]);
1099        store.put_blob(&blob).expect("source");
1100        let tree = Tree::from_entries(vec![
1101            TreeEntry::file("source.rs", blob.hash(), false).expect("entry"),
1102        ]);
1103        store.put_tree(&tree).expect("tree");
1104        let state = State::new_snapshot(
1105            tree.hash(),
1106            vec![],
1107            Attribution::human(Principal::new("user", "user@example.test")),
1108        );
1109        let prepared = SourcePack::prepare(
1110            &store,
1111            &state,
1112            root.path(),
1113            SourceBudget {
1114                max_objects: 16,
1115                max_decoded_bytes: 256 * 1024,
1116            },
1117        )
1118        .expect("prepare exact source");
1119        let (remote, _, _) = fixture(false);
1120        let thread = ThreadRef {
1121            spool: Some(SpoolRef {
1122                id: "spool-test".into(),
1123            }),
1124            id: Some(ThreadId { value: vec![1; 32] }),
1125        };
1126        let receipt = remote
1127            .thread(thread.clone())
1128            .publish_source(
1129                &prepared,
1130                &originals(),
1131                PublicationOptions {
1132                    client_operation_id: "source-upload".into(),
1133                    source: EndpointRef {
1134                        public_key: vec![2; 32],
1135                        kind: EndpointKind::Device as i32,
1136                    },
1137                    sharing_policy_version: vec![],
1138                    checkpoint: None,
1139                },
1140            )
1141            .await
1142            .expect("one Thread-bound publication");
1143        assert_eq!(receipt.thread, Some(thread.clone()));
1144        assert_eq!(
1145            receipt.revision,
1146            Some(RevisionRef {
1147                spool: thread.spool,
1148                revision: Some(revision_ref::Revision::State(
1149                    api::heddle::api::common::StateId {
1150                        value: state.id().as_bytes().to_vec()
1151                    }
1152                ))
1153            })
1154        );
1155    }
1156}