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