native_completion_tests.rs (5604B)
1 use super::{ 2 tests::{config, png, read_request, upload_request}, 3 *, 4 }; 5 use std::sync::Arc; 6 use tokio::{io::AsyncWriteExt, net::TcpListener, sync::Notify}; 7 8 fn descriptor(transaction: &BlossomUploadTransaction) -> Vec<u8> { 9 serde_json::to_vec( 10 &BlobDescriptor::new( 11 transaction.expected_url().clone(), 12 transaction.request().sha256(), 13 transaction.request().byte_size(), 14 MediaType::parse("image/png").unwrap(), 15 1_900_000_000, 16 ) 17 .unwrap(), 18 ) 19 .unwrap() 20 } 21 22 #[tokio::test] 23 async fn retrieval_failures_preserve_native_upload_context_and_passive_evidence() { 24 for unavailable in [false, true] { 25 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); 26 let origin = format!("http://{}", listener.local_addr().unwrap()); 27 let server = tokio::spawn(async move { 28 let (mut stream, _) = listener.accept().await.unwrap(); 29 assert!( 30 String::from_utf8(read_request(&mut stream).await) 31 .unwrap() 32 .starts_with("GET /") 33 ); 34 let body = if unavailable { Vec::new() } else { png(3, 3) }; 35 let status = if unavailable { 36 "503 Service Unavailable" 37 } else { 38 "200 OK" 39 }; 40 let head = format!( 41 "HTTP/1.1 {status}\r\nContent-Type: image/png\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", 42 body.len() 43 ); 44 stream.write_all(head.as_bytes()).await.unwrap(); 45 stream.write_all(&body).await.unwrap(); 46 stream.shutdown().await.unwrap(); 47 }); 48 let slot = crate::transport::BlossomSlot::new(); 49 slot.configure(config(&origin)).unwrap(); 50 let transaction = slot 51 .prepare_upload(upload_request(&origin, png(2, 3))) 52 .unwrap(); 53 let body = descriptor(&transaction); 54 let error = slot 55 .complete_native_upload( 56 transaction, 57 200, 58 Some("application/json"), 59 None, 60 &body, 61 BlossomCancellation::default(), 62 ) 63 .await 64 .unwrap_err(); 65 server.await.unwrap(); 66 assert!( 67 error.possible_orphan(), 68 "A retrieval failure cannot erase the native upload" 69 ); 70 assert_eq!(error.attempts(), 2); 71 assert_eq!(error.retryable(), unavailable); 72 if unavailable { 73 assert_eq!(error.http_status(), Some(503)); 74 } 75 let evidence = slot.evidence().unwrap(); 76 assert!(evidence.possible_orphan()); 77 assert_eq!(evidence.attempts(), 2); 78 assert_eq!( 79 evidence.last_successful_state(), 80 crate::transport::BlossomEvidenceState::UploadVerified 81 ); 82 assert_eq!(evidence.error_code(), Some(error.code())); 83 } 84 } 85 86 #[tokio::test] 87 async fn cancellation_before_retrieval_preserves_the_completed_native_attempt() { 88 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); 89 let origin = format!("http://{}", listener.local_addr().unwrap()); 90 let slot = crate::transport::BlossomSlot::new(); 91 slot.configure(config(&origin)).unwrap(); 92 let transaction = slot 93 .prepare_upload(upload_request(&origin, png(2, 3))) 94 .unwrap(); 95 let body = descriptor(&transaction); 96 let cancellation = BlossomCancellation::default(); 97 cancellation.cancel(); 98 let error = slot 99 .complete_native_upload( 100 transaction, 101 201, 102 Some("application/json"), 103 None, 104 &body, 105 cancellation, 106 ) 107 .await 108 .unwrap_err(); 109 assert_eq!(error.kind(), BlossomErrorKind::Cancelled); 110 assert!(error.possible_orphan()); 111 assert_eq!(error.attempts(), 1); 112 assert!( 113 tokio::time::timeout(Duration::from_millis(20), listener.accept()) 114 .await 115 .is_err() 116 ); 117 } 118 119 #[tokio::test] 120 async fn cancellation_during_retrieval_preserves_both_attempts() { 121 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); 122 let origin = format!("http://{}", listener.local_addr().unwrap()); 123 let entered = Arc::new(Notify::new()); 124 let server = { 125 let entered = entered.clone(); 126 tokio::spawn(async move { 127 let (mut stream, _) = listener.accept().await.unwrap(); 128 let _ = read_request(&mut stream).await; 129 entered.notify_one(); 130 std::future::pending::<()>().await; 131 }) 132 }; 133 let slot = crate::transport::BlossomSlot::new(); 134 slot.configure(config(&origin)).unwrap(); 135 let transaction = slot 136 .prepare_upload(upload_request(&origin, png(2, 3))) 137 .unwrap(); 138 let body = descriptor(&transaction); 139 let cancellation = BlossomCancellation::default(); 140 let result = slot.complete_native_upload( 141 transaction, 142 200, 143 Some("application/json"), 144 None, 145 &body, 146 cancellation.clone(), 147 ); 148 let cancel = async { 149 tokio::time::timeout(Duration::from_secs(5), entered.notified()) 150 .await 151 .unwrap(); 152 cancellation.cancel(); 153 }; 154 let (result, ()) = tokio::join!(result, cancel); 155 server.abort(); 156 assert!(server.await.unwrap_err().is_cancelled()); 157 let error = result.unwrap_err(); 158 assert_eq!(error.kind(), BlossomErrorKind::Cancelled); 159 assert!(error.possible_orphan()); 160 assert_eq!(error.attempts(), 2); 161 }