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};
29
30/// Exact original creator wrappers and signed source operations, including
31/// every foreign integration dependency. The receiver verifies their authority
32/// and complete causal/source closure before any replica admission.
33#[derive(Clone)]
34pub struct PublicationOriginals {
35    pub geneses: Vec<ThreadGenesisRecord>,
36    pub operations: Vec<ReplicationOperations>,
37}
38impl PublicationOriginals {
39    fn validate_bounds(&self) -> Result<(), Error> {
40        if self.geneses.is_empty()
41            || self.geneses.len() > 128
42            || self.operations.is_empty()
43            || self.operations.len() > 10_000
44            || self
45                .operations
46                .iter()
47                .map(|batch| batch.operations.len())
48                .sum::<usize>()
49                > 10_000
50            || self.operations.iter().any(|batch| {
51                batch.operations.is_empty()
52                    || batch.authority_admissions.len() > batch.operations.len()
53            })
54        {
55            return Err(Error::Invalid(
56                "bounded original genesis and source operations required",
57            ));
58        }
59        let mut acceptances = std::collections::BTreeSet::new();
60        for record in self
61            .geneses
62            .iter()
63            .flat_map(|value| &value.boundary_acceptances)
64            .chain(
65                self.operations
66                    .iter()
67                    .flat_map(|value| &value.boundary_acceptances),
68            )
69        {
70            if record.canonical_record.len() > 96 * 1024
71                || record.signatures.len() != 1
72                || record.signatures[0].signature.len() != 64
73                || record.signatures[0].public_key.len() != 32
74            {
75                return Err(Error::Invalid("boundary evidence shape exceeds bounds"));
76            }
77            acceptances.insert(record.canonical_record.as_slice());
78            if acceptances.len() > 128 {
79                return Err(Error::Invalid("boundary acceptance count exceeded"));
80            }
81        }
82        let mut bytes = 0usize;
83        for length in self
84            .geneses
85            .iter()
86            .map(Message::encoded_len)
87            .chain(self.operations.iter().map(Message::encoded_len))
88        {
89            bytes = bytes
90                .checked_add(length)
91                .ok_or(Error::Invalid("publication original size overflow"))?;
92            if length > 256 * 1024 || bytes > 16 * 1024 * 1024 {
93                return Err(Error::Invalid(
94                    "publication original metadata budget exceeded",
95                ));
96            }
97        }
98        if self.geneses.iter().any(|genesis| {
99            genesis.genesis.as_ref().is_none_or(|record| {
100                record.canonical_record.is_empty() || record.signatures.is_empty()
101            })
102        }) || self
103            .operations
104            .iter()
105            .flat_map(|batch| {
106                batch
107                    .operations
108                    .iter()
109                    .chain(&batch.authority_admissions)
110                    .chain(&batch.boundary_acceptances)
111            })
112            .any(|record| record.canonical_record.is_empty() || record.signatures.is_empty())
113        {
114            return Err(Error::Invalid(
115                "original signatures and canonical records required",
116            ));
117        }
118        Ok(())
119    }
120}
121
122#[derive(Debug, thiserror::Error)]
123pub enum Error {
124    #[error(transparent)]
125    Client(#[from] ClientError<transport::Error>),
126    #[error(transparent)]
127    Transport(#[from] transport::Error),
128    #[error("publication source: {0}")]
129    Source(#[from] std::io::Error),
130    #[error("invalid publication: {0}")]
131    Invalid(&'static str),
132    #[error(
133        "publication {limit_name} limit is {limit}, required {actual} for {operations} operations"
134    )]
135    OriginalBudgetExceeded {
136        operations: usize,
137        limit_name: &'static str,
138        limit: usize,
139        actual: usize,
140    },
141    #[error(
142        "publication operation {operation} requires {bytes} encoded bytes; batch limit {batch_limit} bytes, frame limit {frame_limit} bytes"
143    )]
144    OriginalOperationTooLarge {
145        operation: usize,
146        bytes: usize,
147        batch_limit: usize,
148        frame_limit: usize,
149    },
150}
151
152impl<T: RpcTransport<Error = transport::Error>> Remote<T> {
153    /// Publish original source proofs and exact source in one exchange. Inputs are the native pack
154    /// and its index in opening order; neither is buffered in full. The caller
155    /// retains the operation ID and opening to retry after an interrupted call.
156    /// A receipt is returned only after verifying its exact publication scope.
157    pub async fn publish_content<R: AsyncRead + Unpin + Send>(
158        &self,
159        opening: &PublishContentClientFrame,
160        originals: &PublicationOriginals,
161        mut artifacts: [R; 2],
162    ) -> Result<PublicationReceipt, Error> {
163        // Re-batching rebuilds each unit, so refuse HYBRID import authority
164        // here rather than relay originals with it silently dropped (api#307).
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        if open.destination != self.description.endpoint
177            || open.thread.is_none()
178            || open.revision.is_none()
179            || (!open.sharing_policy_version.is_empty() && open.sharing_policy_version.len() != 32)
180            || open.packs.len() != 2
181            || open.packs[0].kind != pack_extent::Kind::NativePack as i32
182            || open.packs[1].kind != pack_extent::Kind::NativeIndex as i32
183        {
184            return Err(Error::Invalid(
185                "endpoint, capture and ordered native artifacts required",
186            ));
187        }
188        let inventory = inventory_digest(&open.packs)?;
189        let mut logical = opening.clone();
190        if let Some(publish_content_client_frame::Body::Open(open)) = logical.body.as_mut() {
191            open.checkpoint = None;
192        }
193        let digest = typed_digest("thread-source-transfer-v1", &logical.encode_to_vec());
194        let expected = TransferCheckpoint {
195            transfer_id: digest[..16].to_vec(),
196            plan_digest: digest.to_vec(),
197            resume_token: Vec::new(),
198            committed_bytes: 0,
199        };
200        let (mut sender, mut responses) = self
201            .api
202            .exchange::<rpc::SyncServicePublishContent>(opening)
203            .await?;
204        let first = responses
205            .next()
206            .await?
207            .ok_or(Error::Invalid("publication ended before admission"))?;
208        let ready = match first.body {
209            Some(publish_content_server_frame::Body::Receipt(receipt)) => {
210                return validate_receipt(receipt, opening, &inventory);
211            }
212            Some(publish_content_server_frame::Body::Ready(ready)) => ready,
213            _ => return Err(Error::Invalid("Ready or replay receipt required")),
214        };
215        crate::hybrid::transfer_ready(&ready).map_err(Error::Invalid)?;
216        if ready.endpoint != open.destination
217            || ready.thread != open.thread
218            || ready.current != open.revision
219            || ready.checkpoint.as_ref() != Some(&expected)
220        {
221            return Err(Error::Invalid("admission differs from publication plan"));
222        }
223        let budget = ready
224            .budget
225            .ok_or(Error::Invalid("publication frame budget required"))?;
226        let frame_limit = budget.max_frame_bytes as usize;
227        if !(1024..=512 * 1024).contains(&frame_limit) {
228            return Err(Error::Invalid("unsupported publication frame budget"));
229        }
230        let originals =
231            batching::bounded_originals(&originals, &opening.client_operation_id, frame_limit)?;
232        originals.validate_bounds()?;
233        let upload = async {
234            for body in originals
235                .geneses
236                .iter()
237                .cloned()
238                .map(publish_content_client_frame::Body::ThreadGenesis)
239                .chain(
240                    originals
241                        .operations
242                        .iter()
243                        .cloned()
244                        .map(publish_content_client_frame::Body::Operations),
245                )
246            {
247                let frame = PublishContentClientFrame {
248                    client_operation_id: opening.client_operation_id.clone(),
249                    body: Some(body),
250                };
251                if frame.encoded_len() > frame_limit {
252                    return Err(Error::Invalid(
253                        "original exceeds negotiated publication frame budget",
254                    ));
255                }
256                sender.send(&frame).await?;
257            }
258            for (artifact, planned) in artifacts.iter_mut().zip(&open.packs) {
259                let mut offset = 0;
260                let mut digest = blake3::Hasher::new();
261                let mut buffer = vec![0; frame_limit / 2];
262                while offset < planned.length {
263                    let length = (planned.length - offset).min(buffer.len() as u64) as usize;
264                    let read = artifact.read(&mut buffer[..length]).await?;
265                    if read == 0 {
266                        return Err(Error::Invalid("artifact ended before declared length"));
267                    }
268                    let data = &buffer[..read];
269                    digest.update(data);
270                    let chunk = PublishContentClientFrame {
271                        client_operation_id: opening.client_operation_id.clone(),
272                        body: Some(publish_content_client_frame::Body::Pack(PackChunk {
273                            extent: Some(PackExtent {
274                                pack: planned.pack.clone(),
275                                kind: planned.kind,
276                                offset,
277                                length: read as u64,
278                                extent_digest: Some(ObjectAddress {
279                                    algorithm: "blake3".into(),
280                                    digest: blake3::hash(data).as_bytes().to_vec(),
281                                }),
282                            }),
283                            data: data.to_vec(),
284                        })),
285                    };
286                    if chunk.encoded_len() > frame_limit {
287                        return Err(Error::Invalid("chunk exceeds negotiated frame budget"));
288                    }
289                    sender.send(&chunk).await?;
290                    offset += read as u64;
291                }
292                if artifact.read(&mut buffer[..1]).await? != 0
293                    || planned
294                        .pack
295                        .as_ref()
296                        .is_none_or(|address| address.digest != digest.finalize().as_bytes())
297                {
298                    return Err(Error::Invalid(
299                        "artifact differs from declared length or digest",
300                    ));
301                }
302            }
303            sender
304                .send(&PublishContentClientFrame {
305                    client_operation_id: opening.client_operation_id.clone(),
306                    body: Some(publish_content_client_frame::Body::Finish(
307                        PublishContentFinish {
308                            checkpoint: Some(expected.clone()),
309                        },
310                    )),
311                })
312                .await?;
313            sender.finish().await?;
314            Ok::<_, Error>(())
315        };
316        let receive = async {
317            while let Some(frame) = responses.next().await? {
318                match frame.body {
319                    Some(publish_content_server_frame::Body::Checkpoint(checkpoint))
320                        if checkpoint == expected => {}
321                    Some(publish_content_server_frame::Body::Receipt(receipt)) => {
322                        return validate_receipt(receipt, opening, &inventory);
323                    }
324                    _ => return Err(Error::Invalid("unexpected publication response")),
325                }
326            }
327            Err(Error::Invalid(
328                "publication ended without a durable receipt",
329            ))
330        };
331        let (_, receipt) = tokio::try_join!(upload, receive)?;
332        Ok(receipt)
333    }
334}
335
336fn validate_receipt(
337    receipt: PublicationReceipt,
338    opening: &PublishContentClientFrame,
339    inventory: &[u8; 32],
340) -> Result<PublicationReceipt, Error> {
341    let Some(publish_content_client_frame::Body::Open(open)) = &opening.body else {
342        return Err(Error::Invalid("Open required"));
343    };
344    crate::hybrid::publication_receipt(&receipt).map_err(Error::Invalid)?;
345    if receipt.client_operation_id != opening.client_operation_id
346        || receipt.destination != open.destination
347        || receipt.thread != open.thread
348        || receipt.revision != open.revision
349        || receipt.sharing_policy_version.len() != 32
350        || (!open.sharing_policy_version.is_empty()
351            && receipt.sharing_policy_version != open.sharing_policy_version)
352    {
353        return Err(Error::Invalid("receipt differs from requested publication"));
354    }
355    match &receipt.outcome {
356        Some(publication_receipt::Outcome::Accepted(_))
357            if receipt.accepted_inventory.as_ref().is_some_and(|address| {
358                address.algorithm == "blake3" && address.digest == inventory
359            }) =>
360        {
361            Ok(receipt)
362        }
363        Some(publication_receipt::Outcome::Rejected(error)) => {
364            Err(transport::Error::Remote(error.clone().into()).into())
365        }
366        _ => Err(Error::Invalid(
367            "receipt did not accept the complete inventory",
368        )),
369    }
370}
371
372/// Exact ordered complete artifact inventory shared by preparation, intent,
373/// upload, and receipt verification. No pack body is buffered here.
374pub fn inventory_digest(packs: &[PackExtent]) -> Result<[u8; 32], Error> {
375    if packs.len() != 2
376        || packs[0].kind != pack_extent::Kind::NativePack as i32
377        || packs[1].kind != pack_extent::Kind::NativeIndex as i32
378    {
379        return Err(Error::Invalid("ordered native pack and index required"));
380    }
381    let mut inventory = Vec::new();
382    let mut total = 0_u64;
383    for extent in packs {
384        let address = extent
385            .pack
386            .as_ref()
387            .ok_or(Error::Invalid("artifact address required"))?;
388        total = total
389            .checked_add(extent.length)
390            .ok_or(Error::Invalid("artifact length overflow"))?;
391        if address.algorithm != "blake3"
392            || address.digest.len() != 32
393            || extent.length == 0
394            || extent.offset != 0
395            || extent.extent_digest.as_ref() != Some(address)
396            || total > 256 * 1024 * 1024
397        {
398            return Err(Error::Invalid("complete bounded BLAKE3 artifacts required"));
399        }
400        extent
401            .encode_length_delimited(&mut inventory)
402            .map_err(|_| Error::Invalid("inventory encoding failed"))?;
403    }
404    Ok(typed_digest("thread-source-inventory-v1", &inventory))
405}
406
407fn typed_digest(kind: &str, bytes: &[u8]) -> [u8; 32] {
408    let mut hasher = blake3::Hasher::new();
409    hasher.update(kind.as_bytes());
410    hasher.update(&(bytes.len() as u64).to_le_bytes());
411    hasher.update(&[0]);
412    hasher.update(bytes);
413    *hasher.finalize().as_bytes()
414}
415
416#[cfg(test)]
417mod tests {
418    use api::v2::{
419        MethodDescriptor,
420        client::{MessageReader, MessageWriter},
421    };
422    use tokio::sync::mpsc;
423
424    use super::*;
425
426    pub(super) struct Reader(mpsc::Receiver<Vec<u8>>);
427    pub(super) struct Writer(Option<mpsc::Sender<Vec<u8>>>);
428    impl MessageReader for Reader {
429        type Error = transport::Error;
430        async fn next(&mut self) -> Result<Option<Vec<u8>>, Self::Error> {
431            Ok(self.0.recv().await)
432        }
433        fn cancel(&mut self) {
434            self.0.close();
435        }
436    }
437    impl MessageWriter for Writer {
438        type Error = transport::Error;
439        async fn send(&mut self, bytes: Vec<u8>) -> Result<(), Self::Error> {
440            self.0
441                .as_ref()
442                .ok_or(transport::Error::Protocol("closed"))?
443                .send(bytes)
444                .await
445                .map_err(|_| transport::Error::Protocol("peer closed"))
446        }
447        async fn finish(&mut self) -> Result<(), Self::Error> {
448            self.0.take();
449            Ok(())
450        }
451        fn abort(&mut self) {
452            self.0.take();
453        }
454    }
455    pub(super) struct Peer {
456        wrong_receipt: bool,
457        original_count: usize,
458    }
459    impl RpcTransport for Peer {
460        type Error = transport::Error;
461        type Reader = Reader;
462        type Writer = Writer;
463        async fn unary(
464            &self,
465            _: &'static MethodDescriptor,
466            _: Vec<u8>,
467        ) -> Result<Vec<u8>, Self::Error> {
468            unreachable!("test exercises only publication");
469        }
470        async fn observe(
471            &self,
472            _: &'static MethodDescriptor,
473            _: Vec<u8>,
474        ) -> Result<Reader, Self::Error> {
475            unreachable!("test exercises only publication");
476        }
477        async fn exchange(
478            &self,
479            _: &'static MethodDescriptor,
480            bytes: Vec<u8>,
481        ) -> Result<(Writer, Reader), Self::Error> {
482            let opening = PublishContentClientFrame::decode(bytes.as_slice())?;
483            let Some(publish_content_client_frame::Body::Open(open)) = opening.body.clone() else {
484                panic!("Open");
485            };
486            let mut logical = opening.clone();
487            if let Some(publish_content_client_frame::Body::Open(open)) = logical.body.as_mut() {
488                open.checkpoint = None;
489            }
490            let digest = typed_digest("thread-source-transfer-v1", &logical.encode_to_vec());
491            let checkpoint = TransferCheckpoint {
492                transfer_id: digest[..16].to_vec(),
493                plan_digest: digest.to_vec(),
494                committed_bytes: 0,
495                resume_token: vec![],
496            };
497            let (tx, mut incoming) = mpsc::channel::<Vec<u8>>(1);
498            let (outgoing, rx) = mpsc::channel(1);
499            let wrong_receipt = self.wrong_receipt;
500            let original_count = self.original_count;
501            tokio::spawn(async move {
502                let send = |body| {
503                    let outgoing = &outgoing;
504                    async move {
505                        outgoing
506                            .send(PublishContentServerFrame { body: Some(body) }.encode_to_vec())
507                            .await
508                    }
509                };
510                send(publish_content_server_frame::Body::Ready(TransferReady {
511                    endpoint: open.destination.clone(),
512                    thread: open.thread.clone(),
513                    current: open.revision.clone(),
514                    checkpoint: Some(checkpoint.clone()),
515                    budget: Some(ReadBudget {
516                        max_frame_bytes: 2048,
517                        ..Default::default()
518                    }),
519                    ..Default::default()
520                }))
521                .await
522                .expect("Ready");
523                let mut lengths = [0_u64; 2];
524                let mut original_counts = [0usize; 2];
525                while let Some(bytes) = incoming.recv().await {
526                    assert!(bytes.len() <= 2048, "negotiated frame size");
527                    let frame = PublishContentClientFrame::decode(bytes.as_slice()).expect("frame");
528                    assert_eq!(frame.client_operation_id, opening.client_operation_id);
529                    match frame.body {
530                        Some(publish_content_client_frame::Body::ThreadGenesis(genesis)) => {
531                            assert_eq!(lengths, [0, 0]);
532                            assert!(genesis.genesis.is_some());
533                            original_counts[0] += 1;
534                            send(publish_content_server_frame::Body::Checkpoint(
535                                checkpoint.clone(),
536                            ))
537                            .await
538                            .expect("original checkpoint");
539                        }
540                        Some(publish_content_client_frame::Body::Operations(batch)) => {
541                            assert_eq!(lengths, [0, 0]);
542                            assert!((1..=128).contains(&batch.operations.len()));
543                            assert!(!batch.operations[0].canonical_record.is_empty());
544                            original_counts[1] += batch.operations.len();
545                            send(publish_content_server_frame::Body::Checkpoint(
546                                checkpoint.clone(),
547                            ))
548                            .await
549                            .expect("original checkpoint");
550                        }
551                        Some(publish_content_client_frame::Body::Pack(chunk)) => {
552                            assert_eq!(
553                                original_counts,
554                                [1, original_count],
555                                "original proofs precede source artifacts"
556                            );
557                            let extent = chunk.extent.expect("extent");
558                            let index = if extent.kind == pack_extent::Kind::NativePack as i32 {
559                                0
560                            } else {
561                                1
562                            };
563                            assert_eq!(extent.offset, lengths[index]);
564                            assert_eq!(extent.length, chunk.data.len() as u64);
565                            assert_eq!(
566                                extent.extent_digest.expect("chunk digest").digest,
567                                blake3::hash(&chunk.data).as_bytes()
568                            );
569                            lengths[index] += extent.length;
570                            // Capacity one in BOTH directions. An upload that
571                            // waits to read responses until all sends complete
572                            // deadlocks after a few chunks here.
573                            send(publish_content_server_frame::Body::Checkpoint(
574                                checkpoint.clone(),
575                            ))
576                            .await
577                            .expect("checkpoint");
578                        }
579                        Some(publish_content_client_frame::Body::Finish(finish)) => {
580                            assert_eq!(finish.checkpoint, Some(checkpoint.clone()));
581                            assert_eq!(lengths, [open.packs[0].length, open.packs[1].length]);
582                            let mut inventory = Vec::new();
583                            for extent in &open.packs {
584                                extent
585                                    .encode_length_delimited(&mut inventory)
586                                    .expect("inventory");
587                            }
588                            let mut receipt = PublicationReceipt {
589                                client_operation_id: opening.client_operation_id.clone(),
590                                destination: open.destination.clone(),
591                                thread: open.thread.clone(),
592                                revision: open.revision.clone(),
593                                sharing_policy_version: if open.sharing_policy_version.is_empty() {
594                                    vec![3; 32]
595                                } else {
596                                    open.sharing_policy_version.clone()
597                                },
598                                accepted_inventory: Some(ObjectAddress {
599                                    algorithm: "blake3".into(),
600                                    digest: typed_digest("thread-source-inventory-v1", &inventory)
601                                        .to_vec(),
602                                }),
603                                outcome: Some(publication_receipt::Outcome::Accepted(
604                                    Applied::default(),
605                                )),
606                                import_authority: None,
607                            };
608                            if wrong_receipt {
609                                receipt.thread = None;
610                            }
611                            let _ =
612                                send(publish_content_server_frame::Body::Receipt(receipt)).await;
613                            break;
614                        }
615                        _ => panic!("unexpected frame"),
616                    }
617                }
618            });
619            Ok((Writer(Some(tx)), Reader(rx)))
620        }
621    }
622
623    pub(super) fn fixture(
624        wrong_receipt: bool,
625    ) -> (
626        Remote<Peer>,
627        PublishContentClientFrame,
628        [std::io::Cursor<Vec<u8>>; 2],
629    ) {
630        let endpoint = EndpointRef {
631            public_key: vec![8; 32],
632            kind: EndpointKind::Weft as i32,
633        };
634        let remote = Remote {
635            api: api::v2::client::Client::new(
636                Peer {
637                    wrong_receipt,
638                    original_count: 1,
639                },
640                ["/heddle.api.v1alpha2.SyncService/PublishContent".into()],
641            ),
642            description: DescribeEndpointResponse {
643                endpoint: Some(endpoint.clone()),
644                ..Default::default()
645            },
646        };
647        let artifacts = [vec![1; 128 * 1024], vec![2; 16 * 1024]];
648        let packs = artifacts
649            .iter()
650            .zip([
651                pack_extent::Kind::NativePack,
652                pack_extent::Kind::NativeIndex,
653            ])
654            .map(|(data, kind)| {
655                let address = ObjectAddress {
656                    algorithm: "blake3".into(),
657                    digest: blake3::hash(data).as_bytes().to_vec(),
658                };
659                PackExtent {
660                    pack: Some(address.clone()),
661                    kind: kind as i32,
662                    offset: 0,
663                    length: data.len() as u64,
664                    extent_digest: Some(address),
665                }
666            })
667            .collect();
668        let open = PublishContentClientFrame {
669            client_operation_id: "op-test".into(),
670            body: Some(publish_content_client_frame::Body::Open(
671                PublishContentOpen {
672                    destination: Some(endpoint),
673                    thread: Some(ThreadRef::default()),
674                    revision: Some(RevisionRef::default()),
675                    packs,
676                    ..Default::default()
677                },
678            )),
679        };
680        (remote, open, artifacts.map(std::io::Cursor::new))
681    }
682    // These byte fixtures exercise framing/backpressure, not original authority
683    // admission; real Iroh tests independently verify canonical signatures.
684    fn originals() -> PublicationOriginals {
685        let record = SignedRecord {
686            format: "transport-fixture".into(),
687            canonical_record: vec![1],
688            signatures: vec![RecordSignature {
689                public_key: vec![2; 32],
690                signature: vec![3; 64],
691            }],
692        };
693        PublicationOriginals {
694            geneses: vec![ThreadGenesisRecord {
695                boundary_acceptances: Vec::new(),
696                genesis: Some(record.clone()),
697                ..Default::default()
698            }],
699            operations: vec![ReplicationOperations {
700                boundary_acceptances: Vec::new(),
701                operations: vec![record],
702                authority_admissions: vec![],
703                import_authority: None,
704            }],
705        }
706    }
707
708    #[test]
709    fn publication_originals_require_bounded_complete_metadata() {
710        let valid = originals();
711        assert!(valid.validate_bounds().is_ok());
712        let mut missing = valid.clone();
713        missing.operations.clear();
714        assert!(matches!(
715            missing.validate_bounds(),
716            Err(Error::Invalid(
717                "bounded original genesis and source operations required"
718            ))
719        ));
720        let mut unsigned = valid.clone();
721        unsigned.operations[0].operations[0].signatures.clear();
722        assert!(matches!(
723            unsigned.validate_bounds(),
724            Err(Error::Invalid(
725                "original signatures and canonical records required"
726            ))
727        ));
728        let mut oversized = valid.clone();
729        oversized.operations[0].operations[0].canonical_record = vec![1; 256 * 1024];
730        assert!(matches!(
731            oversized.validate_bounds(),
732            Err(Error::Invalid(
733                "publication original metadata budget exceeded"
734            ))
735        ));
736        let mut too_many = valid;
737        too_many.geneses = vec![too_many.geneses[0].clone(); 129];
738        assert!(matches!(
739            too_many.validate_bounds(),
740            Err(Error::Invalid(
741                "bounded original genesis and source operations required"
742            ))
743        ));
744    }
745
746    #[tokio::test]
747    async fn publication_drains_checkpoints_while_uploading_under_backpressure() {
748        let (remote, open, artifacts) = fixture(false);
749        let receipt = tokio::time::timeout(
750            std::time::Duration::from_secs(2),
751            remote.publish_content(&open, &originals(), artifacts),
752        )
753        .await
754        .expect("must not deadlock")
755        .expect("publication");
756        assert_eq!(receipt.client_operation_id, "op-test");
757    }
758    #[tokio::test]
759    async fn publication_rebatches_after_ready_before_upload() {
760        let (mut remote, open, artifacts) = fixture(false);
761        remote.api = api::v2::client::Client::new(
762            Peer {
763                wrong_receipt: false,
764                original_count: 140,
765            },
766            ["/heddle.api.v1alpha2.SyncService/PublishContent".into()],
767        );
768        let mut originals = originals();
769        let mut record = originals.operations[0].operations[0].clone();
770        record.canonical_record = vec![1; 300];
771        originals.operations[0].operations = vec![record; 140];
772        let receipt = tokio::time::timeout(
773            std::time::Duration::from_secs(5),
774            remote.publish_content(&open, &originals, artifacts),
775        )
776        .await
777        .expect("must not deadlock")
778        .expect("negotiated bounded publication");
779        assert_eq!(receipt.client_operation_id, "op-test");
780    }
781
782    #[tokio::test]
783    async fn publication_refuses_originals_carrying_hybrid_import_authority() {
784        let (remote, open, artifacts) = fixture(false);
785        let mut originals = originals();
786        originals.operations[0].import_authority = Some(ImportPublicProofBundleV1::default());
787        let error = remote
788            .publish_content(&open, &originals, artifacts)
789            .await
790            .expect_err("HYBRID originals must not be relayed with the bundle dropped");
791        assert!(
792            matches!(error, Error::Invalid(message) if message.contains("api#307")),
793            "{error}"
794        );
795    }
796
797    #[tokio::test]
798    async fn publication_refuses_a_receipt_for_another_scope() {
799        let (remote, open, artifacts) = fixture(true);
800        let error = remote
801            .publish_content(&open, &originals(), artifacts)
802            .await
803            .expect_err("mismatched scope");
804        assert!(matches!(
805            error,
806            Error::Invalid("receipt differs from requested publication")
807        ));
808    }
809    #[test]
810    fn publication_policy_is_optional_cas_and_receipt_reports_actual_frontier() {
811        let (_, mut opening, _) = fixture(false);
812        let inventory = [7; 32];
813        let Some(publish_content_client_frame::Body::Open(open)) = &opening.body else {
814            panic!("opening")
815        };
816        let receipt = PublicationReceipt {
817            client_operation_id: opening.client_operation_id.clone(),
818            destination: open.destination.clone(),
819            thread: open.thread.clone(),
820            revision: open.revision.clone(),
821            sharing_policy_version: vec![3; 32],
822            accepted_inventory: Some(ObjectAddress {
823                algorithm: "blake3".into(),
824                digest: inventory.to_vec(),
825            }),
826            outcome: Some(publication_receipt::Outcome::Accepted(Applied::default())),
827            import_authority: None,
828        };
829        validate_receipt(receipt.clone(), &opening, &inventory)
830            .expect("one-shot upload returns actual policy without CAS");
831        let mut hybrid = receipt.clone();
832        hybrid.import_authority = Some(ImportPublicProofBundleV1::default());
833        assert!(
834            matches!(
835                validate_receipt(hybrid, &opening, &inventory),
836                Err(Error::Invalid(message)) if message.contains("api#307")
837            ),
838            "a receipt carrying HYBRID import authority must be refused, never ignored"
839        );
840        let mut absent = receipt.clone();
841        absent.sharing_policy_version.clear();
842        assert!(
843            validate_receipt(absent, &opening, &inventory).is_err(),
844            "receipt must report actual policy frontier"
845        );
846        let Some(publish_content_client_frame::Body::Open(open)) = &mut opening.body else {
847            panic!("opening")
848        };
849        open.sharing_policy_version = vec![4; 32];
850        assert!(
851            validate_receipt(receipt.clone(), &opening, &inventory).is_err(),
852            "explicit policy CAS cannot silently accept another version"
853        );
854        let Some(publish_content_client_frame::Body::Open(open)) = &mut opening.body else {
855            panic!("opening")
856        };
857        open.sharing_policy_version = vec![3; 32];
858        validate_receipt(receipt, &opening, &inventory).expect("matching explicit policy version");
859    }
860
861    #[cfg(feature = "replication")]
862    #[test]
863    fn publication_digests_match_the_native_typed_hash_format() {
864        assert_eq!(
865            typed_digest("thread-source-inventory-v1", b"bytes"),
866            *heddle_object_model::object::ContentHash::compute_typed(
867                "thread-source-inventory-v1",
868                b"bytes"
869            )
870            .as_bytes()
871        );
872    }
873
874    #[cfg(feature = "source-transfer")]
875    #[tokio::test]
876    async fn thread_publication_prepares_only_selected_source_and_binds_its_revision() {
877        use objects::{
878            object::{Attribution, Blob, Principal, State, Tree, TreeEntry},
879            store::{FsStore, ObjectStore},
880        };
881        let root = tempfile::tempdir().expect("local source scratch");
882        let store = FsStore::new(root.path().join("objects"));
883        store.init().expect("store");
884        let blob = Blob::new(vec![1; 128 * 1024]);
885        store.put_blob(&blob).expect("source");
886        let tree = Tree::from_entries(vec![
887            TreeEntry::file("source.rs", blob.hash(), false).expect("entry"),
888        ]);
889        store.put_tree(&tree).expect("tree");
890        let state = State::new_snapshot(
891            tree.hash(),
892            vec![],
893            Attribution::human(Principal::new("user", "user@example.test")),
894        );
895        let prepared = SourcePack::prepare(
896            &store,
897            &state,
898            root.path(),
899            SourceBudget {
900                max_objects: 16,
901                max_decoded_bytes: 256 * 1024,
902            },
903        )
904        .expect("prepare exact source");
905        let (remote, _, _) = fixture(false);
906        let thread = ThreadRef {
907            spool: Some(SpoolRef {
908                id: "spool-test".into(),
909            }),
910            id: Some(ThreadId { value: vec![1; 32] }),
911        };
912        let receipt = remote
913            .thread(thread.clone())
914            .publish_source(
915                &prepared,
916                &originals(),
917                PublicationOptions {
918                    client_operation_id: "source-upload".into(),
919                    source: EndpointRef {
920                        public_key: vec![2; 32],
921                        kind: EndpointKind::Device as i32,
922                    },
923                    sharing_policy_version: vec![],
924                    checkpoint: None,
925                },
926            )
927            .await
928            .expect("one Thread-bound publication");
929        assert_eq!(receipt.thread, Some(thread.clone()));
930        assert_eq!(
931            receipt.revision,
932            Some(RevisionRef {
933                spool: thread.spool,
934                revision: Some(revision_ref::Revision::State(
935                    api::heddle::api::common::StateId {
936                        value: state.id().as_bytes().to_vec()
937                    }
938                ))
939            })
940        );
941    }
942}