commit bb54297d208f8e3817485d1c1a63a386f457349a
parent ac392da942896a6b676f063929af16971e201b0e
Author: triesap <tyson@radroots.org>
Date: Sat, 12 Sep 2026 20:34:00 +0000
blossom: preserve native upload context on retrieval failure
- Retain possible remote effects after accepted native completion
- Count the native attempt alongside bounded retrieval attempts
- Cover failed retrieval and cancellation with loopback regressions
- Pass unchanged API coverage generator and workspace gates
Diffstat:
4 files changed, 183 insertions(+), 14 deletions(-)
diff --git a/crates/sdk/src/adapters/blossom.rs b/crates/sdk/src/adapters/blossom.rs
@@ -374,17 +374,19 @@ pub(crate) async fn complete_native_upload(
)
})?;
let verified_upload = verify_descriptor(&request, &expected_url, descriptor, 1)?;
- let retrieved = retrieve(
- config,
- BlossomInboundRequest::new(
- verified_upload.url().as_blob_url().clone(),
- Some(request.media_type().clone()),
- Some(request.byte_size()),
- Some(request.dimensions()),
- )?,
- cancellation,
+ let inbound = BlossomInboundRequest::new(
+ verified_upload.url().as_blob_url().clone(),
+ Some(request.media_type().clone()),
+ Some(request.byte_size()),
+ Some(request.dimensions()),
)
- .await?;
+ .map_err(|error| with_operation(error, true, 1))?;
+ let retrieved = retrieve(config, inbound, cancellation)
+ .await
+ .map_err(|error| {
+ let attempts = 1_u8.saturating_add(error.attempts());
+ with_operation(error, true, attempts)
+ })?;
if retrieved.bytes() != request.bytes() {
return Err(failure(
BlossomErrorKind::RetrievedBytesMismatch,
@@ -1459,6 +1461,9 @@ fn little_u24(bytes: &[u8]) -> Option<u32> {
}
#[cfg(test)]
+mod native_completion_tests;
+
+#[cfg(test)]
mod tests {
use super::*;
use radroots_blossom::Sha256;
@@ -1468,7 +1473,7 @@ mod tests {
static LOOPBACK_TEST_GUARD: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
- fn png(width: u32, height: u32) -> Vec<u8> {
+ pub(super) fn png(width: u32, height: u32) -> Vec<u8> {
let mut bytes = b"\x89PNG\r\n\x1a\n\0\0\0\rIHDR".to_vec();
bytes.extend_from_slice(&width.to_be_bytes());
bytes.extend_from_slice(&height.to_be_bytes());
@@ -1877,7 +1882,7 @@ mod tests {
Stall,
}
- async fn read_request(stream: &mut tokio::net::TcpStream) -> Vec<u8> {
+ pub(super) async fn read_request(stream: &mut tokio::net::TcpStream) -> Vec<u8> {
let mut request = Vec::new();
let header_end = loop {
let mut chunk = [0_u8; 1024];
@@ -2079,7 +2084,7 @@ mod tests {
(origin, task)
}
- fn upload_request(_origin: &str, bytes: Vec<u8>) -> BlossomUploadRequest {
+ pub(super) fn upload_request(_origin: &str, bytes: Vec<u8>) -> BlossomUploadRequest {
BlossomUploadRequest::new(
Arc::from(bytes),
MediaType::parse("image/png").expect("media type"),
@@ -2089,7 +2094,7 @@ mod tests {
.expect("upload request")
}
- fn config(origin: &str) -> BlossomConfig {
+ pub(super) fn config(origin: &str) -> BlossomConfig {
BlossomConfig::from_profile(simulator_profile(origin))
.with_network_policy(
Duration::from_secs(2),
diff --git a/crates/sdk/src/adapters/blossom/native_completion_tests.rs b/crates/sdk/src/adapters/blossom/native_completion_tests.rs
@@ -0,0 +1,161 @@
+use super::{
+ tests::{config, png, read_request, upload_request},
+ *,
+};
+use std::sync::Arc;
+use tokio::{io::AsyncWriteExt, net::TcpListener, sync::Notify};
+
+fn descriptor(transaction: &BlossomUploadTransaction) -> Vec<u8> {
+ serde_json::to_vec(
+ &BlobDescriptor::new(
+ transaction.expected_url().clone(),
+ transaction.request().sha256(),
+ transaction.request().byte_size(),
+ MediaType::parse("image/png").unwrap(),
+ 1_900_000_000,
+ )
+ .unwrap(),
+ )
+ .unwrap()
+}
+
+#[tokio::test]
+async fn retrieval_failures_preserve_native_upload_context_and_passive_evidence() {
+ for unavailable in [false, true] {
+ let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let origin = format!("http://{}", listener.local_addr().unwrap());
+ let server = tokio::spawn(async move {
+ let (mut stream, _) = listener.accept().await.unwrap();
+ assert!(
+ String::from_utf8(read_request(&mut stream).await)
+ .unwrap()
+ .starts_with("GET /")
+ );
+ let body = if unavailable { Vec::new() } else { png(3, 3) };
+ let status = if unavailable {
+ "503 Service Unavailable"
+ } else {
+ "200 OK"
+ };
+ let head = format!(
+ "HTTP/1.1 {status}\r\nContent-Type: image/png\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
+ body.len()
+ );
+ stream.write_all(head.as_bytes()).await.unwrap();
+ stream.write_all(&body).await.unwrap();
+ stream.shutdown().await.unwrap();
+ });
+ let slot = crate::transport::BlossomSlot::new();
+ slot.configure(config(&origin)).unwrap();
+ let transaction = slot
+ .prepare_upload(upload_request(&origin, png(2, 3)))
+ .unwrap();
+ let body = descriptor(&transaction);
+ let error = slot
+ .complete_native_upload(
+ transaction,
+ 200,
+ Some("application/json"),
+ None,
+ &body,
+ BlossomCancellation::default(),
+ )
+ .await
+ .unwrap_err();
+ server.await.unwrap();
+ assert!(
+ error.possible_orphan(),
+ "A retrieval failure cannot erase the native upload"
+ );
+ assert_eq!(error.attempts(), 2);
+ assert_eq!(error.retryable(), unavailable);
+ if unavailable {
+ assert_eq!(error.http_status(), Some(503));
+ }
+ let evidence = slot.evidence().unwrap();
+ assert!(evidence.possible_orphan());
+ assert_eq!(evidence.attempts(), 2);
+ assert_eq!(
+ evidence.last_successful_state(),
+ crate::transport::BlossomEvidenceState::UploadVerified
+ );
+ assert_eq!(evidence.error_code(), Some(error.code()));
+ }
+}
+
+#[tokio::test]
+async fn cancellation_before_retrieval_preserves_the_completed_native_attempt() {
+ let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let origin = format!("http://{}", listener.local_addr().unwrap());
+ let slot = crate::transport::BlossomSlot::new();
+ slot.configure(config(&origin)).unwrap();
+ let transaction = slot
+ .prepare_upload(upload_request(&origin, png(2, 3)))
+ .unwrap();
+ let body = descriptor(&transaction);
+ let cancellation = BlossomCancellation::default();
+ cancellation.cancel();
+ let error = slot
+ .complete_native_upload(
+ transaction,
+ 201,
+ Some("application/json"),
+ None,
+ &body,
+ cancellation,
+ )
+ .await
+ .unwrap_err();
+ assert_eq!(error.kind(), BlossomErrorKind::Cancelled);
+ assert!(error.possible_orphan());
+ assert_eq!(error.attempts(), 1);
+ assert!(
+ tokio::time::timeout(Duration::from_millis(20), listener.accept())
+ .await
+ .is_err()
+ );
+}
+
+#[tokio::test]
+async fn cancellation_during_retrieval_preserves_both_attempts() {
+ let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let origin = format!("http://{}", listener.local_addr().unwrap());
+ let entered = Arc::new(Notify::new());
+ let server = {
+ let entered = entered.clone();
+ tokio::spawn(async move {
+ let (mut stream, _) = listener.accept().await.unwrap();
+ let _ = read_request(&mut stream).await;
+ entered.notify_one();
+ std::future::pending::<()>().await;
+ })
+ };
+ let slot = crate::transport::BlossomSlot::new();
+ slot.configure(config(&origin)).unwrap();
+ let transaction = slot
+ .prepare_upload(upload_request(&origin, png(2, 3)))
+ .unwrap();
+ let body = descriptor(&transaction);
+ let cancellation = BlossomCancellation::default();
+ let result = slot.complete_native_upload(
+ transaction,
+ 200,
+ Some("application/json"),
+ None,
+ &body,
+ cancellation.clone(),
+ );
+ let cancel = async {
+ tokio::time::timeout(Duration::from_secs(5), entered.notified())
+ .await
+ .unwrap();
+ cancellation.cancel();
+ };
+ let (result, ()) = tokio::join!(result, cancel);
+ server.abort();
+ assert!(server.await.unwrap_err().is_cancelled());
+ let error = result.unwrap_err();
+ assert_eq!(error.kind(), BlossomErrorKind::Cancelled);
+ assert!(error.possible_orphan());
+ assert_eq!(error.attempts(), 2);
+}
diff --git a/crates/sdk/src/transport.rs b/crates/sdk/src/transport.rs
@@ -1455,6 +1455,8 @@ impl BlossomSlot {
/// Verifies a host-executed BUD-02 response, then performs the canonical
/// BUD-01 exact-byte retrieval before returning an upload receipt.
+ /// Later verification failure or cancellation preserves the native upload's
+ /// possible remote effect and includes that upload in the attempt count.
pub async fn complete_native_upload(
&self,
transaction: BlossomUploadTransaction,
diff --git a/crates/sdk/tests/package_boundary.rs b/crates/sdk/tests/package_boundary.rs
@@ -155,6 +155,7 @@ fn package_contains_only_reachable_sources_and_registered_targets() {
BTreeSet::from([
"adapters/mod.rs".to_owned(),
"adapters/blossom.rs".to_owned(),
+ "adapters/blossom/native_completion_tests.rs".to_owned(),
"adapters/radrootsd.rs".to_owned(),
"capability.rs".to_owned(),
"client.rs".to_owned(),