1use axum::{
7 body::{Body, Bytes},
8 extract::{Path, Query, Request, State},
9 http::{header, HeaderMap, Response, StatusCode},
10 middleware::Next,
11 response::IntoResponse,
12 Json,
13};
14use base64::Engine;
15use hashtree_blossom::{
16 batch_upload_hash_list_digest, BatchUploadItem, BlossomClient, BATCH_UPLOAD_HASH_LIST_AUTH_TAG,
17};
18use hashtree_core::from_hex;
19use nostr::Keys;
20use serde::{Deserialize, Serialize};
21use sha2::{Digest, Sha256};
22use std::collections::HashSet;
23use std::sync::{
24 atomic::{AtomicU64, AtomicUsize, Ordering},
25 Arc, Mutex, OnceLock,
26};
27use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
28use tokio::sync::{mpsc, OwnedSemaphorePermit, Semaphore};
29
30use super::auth::AppState;
31use super::blob_read::{
32 run_blob_metadata_read, run_blob_write, BlobIoTaskError, BLOB_READ_BUSY, BLOB_WRITE_BUSY,
33};
34use super::ingest_filter::{
35 content_type_base, is_chk_content_type, validate_untrusted_blob, IngestRejection,
36};
37use super::mime::get_mime_type;
38
39const BLOSSOM_AUTH_KIND: u16 = 24242;
41
42const IMMUTABLE_CACHE_CONTROL: &str = "public, max-age=31536000, immutable";
44const NOT_FOUND_CACHE_CONTROL: &str = "no-store";
45const IMMUTABLE_NOT_FOUND_CACHE_CONTROL: &str = "public, max-age=0, s-maxage=5";
46const OPTIMISTIC_UPLOAD_QUEUE_TIMEOUT_MS_ENV: &str = "HTREE_OPTIMISTIC_UPLOAD_QUEUE_TIMEOUT_MS";
47const DEFAULT_OPTIMISTIC_UPLOAD_QUEUE_TIMEOUT_MS: u64 = 15_000;
48const BLOSSOM_PUBLIC_BASE_URL_ENV: &str = "HTREE_BLOSSOM_PUBLIC_BASE_URL";
49const LEGACY_BLOSSOM_PUBLIC_BASE_URL_ENV: &str = "HASHTREE_BLOSSOM_PUBLIC_BASE_URL";
50
51pub const DEFAULT_MAX_UPLOAD_SIZE: usize = 5 * 1024 * 1024;
53pub const MAX_SINGLE_UPLOAD_BODY_BYTES: usize = 64 * 1024 * 1024;
54const OPTIMISTIC_UPLOAD_MIN_QUEUE_CHARGE_BYTES: usize = 256 * 1024;
55const MAX_BATCH_UPLOAD_BLOBS: usize = 1024;
56pub const MAX_BATCH_UPLOAD_BYTES: usize = 64 * 1024 * 1024;
57pub const MAX_BATCH_UPLOAD_JSON_BODY_BYTES: usize = 96 * 1024 * 1024;
58const BINARY_BATCH_UPLOAD_MAGIC: &[u8; 8] = b"HTBBV1\0\0";
59const MAX_BINARY_BATCH_CONTENT_TYPE_BYTES: usize = 1024;
60pub const MAX_BATCH_UPLOAD_BINARY_BODY_BYTES: usize = MAX_BATCH_UPLOAD_BYTES
61 + BINARY_BATCH_UPLOAD_MAGIC.len()
62 + 4
63 + (MAX_BATCH_UPLOAD_BLOBS * (32 + 2 + MAX_BINARY_BATCH_CONTENT_TYPE_BYTES + 8));
64const MAX_UPLOAD_CHECK_HASHES: usize = 10_000;
65const SLOW_BATCH_UPLOAD_LOG_MS_ENV: &str = "HTREE_SLOW_BATCH_UPLOAD_LOG_MS";
66const BLOSSOM_REPLICA_UPLOAD_CONCURRENCY_ENV: &str = "HTREE_BLOSSOM_REPLICA_UPLOAD_CONCURRENCY";
67const DEFAULT_BLOSSOM_REPLICA_UPLOAD_CONCURRENCY: usize = 4;
68const BLOSSOM_REPLICA_UPLOAD_ATTEMPTS_ENV: &str = "HTREE_BLOSSOM_REPLICA_UPLOAD_ATTEMPTS";
69const DEFAULT_BLOSSOM_REPLICA_UPLOAD_ATTEMPTS: usize = 3;
70const BLOSSOM_REPLICA_COALESCE_MAX_BLOBS_ENV: &str = "HTREE_BLOSSOM_REPLICA_COALESCE_MAX_BLOBS";
71const DEFAULT_BLOSSOM_REPLICA_COALESCE_MAX_BLOBS: usize = 64;
72const BLOSSOM_REPLICA_COALESCE_MAX_BYTES_ENV: &str = "HTREE_BLOSSOM_REPLICA_COALESCE_MAX_BYTES";
73const DEFAULT_BLOSSOM_REPLICA_COALESCE_MAX_BYTES: usize = 16 * 1024 * 1024;
74const BLOSSOM_REPLICA_COALESCE_FLUSH_MS_ENV: &str = "HTREE_BLOSSOM_REPLICA_COALESCE_FLUSH_MS";
75const DEFAULT_BLOSSOM_REPLICA_COALESCE_FLUSH_MS: u64 = 25;
76const BLOSSOM_REPLICA_COALESCE_QUEUE_JOBS_ENV: &str = "HTREE_BLOSSOM_REPLICA_COALESCE_QUEUE_JOBS";
77const DEFAULT_BLOSSOM_REPLICA_COALESCE_QUEUE_JOBS: usize = 1024;
78
79fn slow_batch_upload_log_ms() -> Option<u128> {
80 std::env::var(SLOW_BATCH_UPLOAD_LOG_MS_ENV)
81 .ok()
82 .and_then(|value| value.parse::<u128>().ok())
83 .filter(|value| *value > 0)
84}
85
86fn blossom_replica_upload_concurrency() -> usize {
87 std::env::var(BLOSSOM_REPLICA_UPLOAD_CONCURRENCY_ENV)
88 .ok()
89 .and_then(|value| value.parse::<usize>().ok())
90 .filter(|value| *value > 0)
91 .unwrap_or(DEFAULT_BLOSSOM_REPLICA_UPLOAD_CONCURRENCY)
92}
93
94fn blossom_replica_upload_attempts() -> usize {
95 std::env::var(BLOSSOM_REPLICA_UPLOAD_ATTEMPTS_ENV)
96 .ok()
97 .and_then(|value| value.parse::<usize>().ok())
98 .filter(|value| *value > 0)
99 .unwrap_or(DEFAULT_BLOSSOM_REPLICA_UPLOAD_ATTEMPTS)
100}
101
102fn blossom_replica_coalesce_max_blobs() -> usize {
103 std::env::var(BLOSSOM_REPLICA_COALESCE_MAX_BLOBS_ENV)
104 .ok()
105 .and_then(|value| value.parse::<usize>().ok())
106 .filter(|value| *value > 0)
107 .unwrap_or(DEFAULT_BLOSSOM_REPLICA_COALESCE_MAX_BLOBS)
108 .min(MAX_BATCH_UPLOAD_BLOBS)
109}
110
111fn blossom_replica_coalesce_max_bytes() -> usize {
112 std::env::var(BLOSSOM_REPLICA_COALESCE_MAX_BYTES_ENV)
113 .ok()
114 .and_then(|value| value.parse::<usize>().ok())
115 .filter(|value| *value > 0)
116 .unwrap_or(DEFAULT_BLOSSOM_REPLICA_COALESCE_MAX_BYTES)
117 .min(MAX_BATCH_UPLOAD_BYTES)
118}
119
120fn blossom_replica_coalesce_flush_delay() -> Duration {
121 let millis = std::env::var(BLOSSOM_REPLICA_COALESCE_FLUSH_MS_ENV)
122 .ok()
123 .and_then(|value| value.parse::<u64>().ok())
124 .unwrap_or(DEFAULT_BLOSSOM_REPLICA_COALESCE_FLUSH_MS);
125 Duration::from_millis(millis)
126}
127
128fn blossom_replica_coalesce_queue_jobs() -> usize {
129 std::env::var(BLOSSOM_REPLICA_COALESCE_QUEUE_JOBS_ENV)
130 .ok()
131 .and_then(|value| value.parse::<usize>().ok())
132 .filter(|value| *value > 0)
133 .unwrap_or(DEFAULT_BLOSSOM_REPLICA_COALESCE_QUEUE_JOBS)
134}
135
136fn blossom_replica_upload_semaphore() -> Arc<Semaphore> {
137 static SEMAPHORE: OnceLock<Arc<Semaphore>> = OnceLock::new();
138 SEMAPHORE
139 .get_or_init(|| Arc::new(Semaphore::new(blossom_replica_upload_concurrency())))
140 .clone()
141}
142
143#[derive(Debug, Clone, Copy)]
144pub(super) struct OptimisticUploadQueueSnapshot {
145 pub enabled: bool,
146 pub max_bytes: usize,
147 pub available_bytes: usize,
148 pub reserved_bytes: usize,
149 pub in_flight: usize,
150 pub queue_timeout_ms: u64,
151}
152
153#[derive(Debug, Clone, Copy)]
154pub(super) struct BlossomUploadReplicaQueueSnapshot {
155 pub enabled: bool,
156 pub target_count: usize,
157 pub max_bytes: usize,
158 pub available_bytes: usize,
159 pub reserved_bytes: usize,
160 pub coalesce_queue_capacity_jobs: usize,
161 pub coalesce_queued_jobs: usize,
162 pub coalesce_max_blobs: usize,
163 pub coalesce_max_bytes: usize,
164 pub coalesce_flush_ms: u64,
165 pub upload_concurrency: usize,
166 pub in_flight_batches: usize,
167 pub accepted_batches: u64,
168 pub accepted_blobs: u64,
169 pub uploaded_blobs: u64,
170 pub replicated_bytes: u64,
171 pub failed_batches: u64,
172 pub skipped_jobs: u64,
173 pub fallback_batches: u64,
174 pub fallback_uploaded_blobs: u64,
175 pub fallback_failed_blobs: u64,
176}
177
178#[derive(Default)]
179struct BlossomUploadReplicaMetrics {
180 coalesce_queued_jobs: AtomicUsize,
181 in_flight_batches: AtomicUsize,
182 accepted_batches: AtomicU64,
183 accepted_blobs: AtomicU64,
184 uploaded_blobs: AtomicU64,
185 replicated_bytes: AtomicU64,
186 failed_batches: AtomicU64,
187 skipped_jobs: AtomicU64,
188 fallback_batches: AtomicU64,
189 fallback_uploaded_blobs: AtomicU64,
190 fallback_failed_blobs: AtomicU64,
191}
192
193fn blossom_upload_replica_metrics() -> &'static BlossomUploadReplicaMetrics {
194 static METRICS: OnceLock<BlossomUploadReplicaMetrics> = OnceLock::new();
195 METRICS.get_or_init(BlossomUploadReplicaMetrics::default)
196}
197
198struct BlossomReplicaInFlightGuard;
199
200impl BlossomReplicaInFlightGuard {
201 fn new() -> Self {
202 blossom_upload_replica_metrics()
203 .in_flight_batches
204 .fetch_add(1, Ordering::Relaxed);
205 Self
206 }
207}
208
209impl Drop for BlossomReplicaInFlightGuard {
210 fn drop(&mut self) {
211 blossom_upload_replica_metrics()
212 .in_flight_batches
213 .fetch_sub(1, Ordering::Relaxed);
214 }
215}
216
217struct PreparedBlossomUploadReplication {
218 servers: Vec<String>,
219 keys: Arc<Keys>,
220 permit: OwnedSemaphorePermit,
221 total_bytes: usize,
222}
223
224struct BlossomReplicaUploadJob {
225 prepared: PreparedBlossomUploadReplication,
226 items: Vec<BatchUploadItem>,
227}
228
229struct BlossomReplicaUploadBatch {
230 servers: Vec<String>,
231 keys: Arc<Keys>,
232 permits: Vec<OwnedSemaphorePermit>,
233 total_bytes: usize,
234 data_bytes: usize,
235 items: Vec<BatchUploadItem>,
236}
237
238impl BlossomReplicaUploadBatch {
239 fn from_job(job: BlossomReplicaUploadJob) -> Self {
240 let PreparedBlossomUploadReplication {
241 servers,
242 keys,
243 permit,
244 total_bytes,
245 } = job.prepared;
246 let data_bytes = job.items.iter().map(|item| item.data.len()).sum();
247 Self {
248 servers,
249 keys,
250 permits: vec![permit],
251 total_bytes,
252 data_bytes,
253 items: job.items,
254 }
255 }
256
257 fn can_append(
258 &self,
259 job: &BlossomReplicaUploadJob,
260 max_blobs: usize,
261 max_bytes: usize,
262 ) -> bool {
263 if self.servers != job.prepared.servers || !Arc::ptr_eq(&self.keys, &job.prepared.keys) {
264 return false;
265 }
266 let job_bytes = job.items.iter().map(|item| item.data.len()).sum::<usize>();
267 self.items.len().saturating_add(job.items.len()) <= max_blobs
268 && self.data_bytes.saturating_add(job_bytes) <= max_bytes
269 }
270
271 fn append(&mut self, job: BlossomReplicaUploadJob) {
272 let PreparedBlossomUploadReplication {
273 permit,
274 total_bytes,
275 ..
276 } = job.prepared;
277 self.permits.push(permit);
278 self.total_bytes = self.total_bytes.saturating_add(total_bytes);
279 self.data_bytes = self
280 .data_bytes
281 .saturating_add(job.items.iter().map(|item| item.data.len()).sum::<usize>());
282 self.items.extend(job.items);
283 }
284
285 fn reached_limits(&self, max_blobs: usize, max_bytes: usize) -> bool {
286 self.items.len() >= max_blobs || self.data_bytes >= max_bytes
287 }
288}
289
290pub struct BlossomUploadReplicaScheduler {
292 sender: Mutex<Option<mpsc::Sender<BlossomReplicaUploadJob>>>,
293}
294
295impl BlossomUploadReplicaScheduler {
296 pub fn new() -> Self {
297 Self {
298 sender: Mutex::new(None),
299 }
300 }
301
302 fn schedule(&self, job: BlossomReplicaUploadJob) -> Result<(), BlossomReplicaUploadJob> {
303 let max_blobs = blossom_replica_coalesce_max_blobs();
304 let flush_delay = blossom_replica_coalesce_flush_delay();
305 if max_blobs <= 1 || flush_delay.is_zero() {
306 return Err(job);
307 }
308
309 let mut job = job;
310 for _ in 0..2 {
311 let sender = self.sender();
312 blossom_upload_replica_metrics()
313 .coalesce_queued_jobs
314 .fetch_add(1, Ordering::Relaxed);
315 match sender.try_send(job) {
316 Ok(()) => return Ok(()),
317 Err(mpsc::error::TrySendError::Full(returned)) => {
318 blossom_upload_replica_metrics()
319 .coalesce_queued_jobs
320 .fetch_sub(1, Ordering::Relaxed);
321 return Err(returned);
322 }
323 Err(mpsc::error::TrySendError::Closed(returned)) => {
324 blossom_upload_replica_metrics()
325 .coalesce_queued_jobs
326 .fetch_sub(1, Ordering::Relaxed);
327 self.clear_sender();
328 job = returned;
329 }
330 }
331 }
332 Err(job)
333 }
334
335 fn sender(&self) -> mpsc::Sender<BlossomReplicaUploadJob> {
336 let mut guard = self.sender.lock().unwrap_or_else(|err| err.into_inner());
337 if guard.as_ref().is_some_and(|sender| sender.is_closed()) {
338 *guard = None;
339 }
340 if let Some(sender) = guard.as_ref() {
341 return sender.clone();
342 }
343
344 let (sender, receiver) = mpsc::channel(blossom_replica_coalesce_queue_jobs());
345 tokio::spawn(blossom_replica_coalescer_worker(receiver));
346 *guard = Some(sender.clone());
347 sender
348 }
349
350 fn clear_sender(&self) {
351 let mut guard = self.sender.lock().unwrap_or_else(|err| err.into_inner());
352 *guard = None;
353 }
354}
355
356impl Default for BlossomUploadReplicaScheduler {
357 fn default() -> Self {
358 Self::new()
359 }
360}
361
362pub(super) fn blossom_upload_replica_queue_snapshot(
363 state: &AppState,
364) -> BlossomUploadReplicaQueueSnapshot {
365 let available_bytes = state.blossom_upload_replica_queue.available_permits();
366 let metrics = blossom_upload_replica_metrics();
367 BlossomUploadReplicaQueueSnapshot {
368 enabled: !state.blossom_upload_replicas.is_empty(),
369 target_count: state.blossom_upload_replicas.len(),
370 max_bytes: state.blossom_upload_replica_queue_bytes,
371 available_bytes,
372 reserved_bytes: state
373 .blossom_upload_replica_queue_bytes
374 .saturating_sub(available_bytes),
375 coalesce_queue_capacity_jobs: blossom_replica_coalesce_queue_jobs(),
376 coalesce_queued_jobs: metrics.coalesce_queued_jobs.load(Ordering::Relaxed),
377 coalesce_max_blobs: blossom_replica_coalesce_max_blobs(),
378 coalesce_max_bytes: blossom_replica_coalesce_max_bytes(),
379 coalesce_flush_ms: duration_millis_u64(blossom_replica_coalesce_flush_delay()),
380 upload_concurrency: blossom_replica_upload_concurrency(),
381 in_flight_batches: metrics.in_flight_batches.load(Ordering::Relaxed),
382 accepted_batches: metrics.accepted_batches.load(Ordering::Relaxed),
383 accepted_blobs: metrics.accepted_blobs.load(Ordering::Relaxed),
384 uploaded_blobs: metrics.uploaded_blobs.load(Ordering::Relaxed),
385 replicated_bytes: metrics.replicated_bytes.load(Ordering::Relaxed),
386 failed_batches: metrics.failed_batches.load(Ordering::Relaxed),
387 skipped_jobs: metrics.skipped_jobs.load(Ordering::Relaxed),
388 fallback_batches: metrics.fallback_batches.load(Ordering::Relaxed),
389 fallback_uploaded_blobs: metrics.fallback_uploaded_blobs.load(Ordering::Relaxed),
390 fallback_failed_blobs: metrics.fallback_failed_blobs.load(Ordering::Relaxed),
391 }
392}
393
394#[allow(clippy::result_large_err)]
397fn check_write_access(state: &AppState, pubkey: &str) -> Result<(), Response<Body>> {
398 if is_allowed_write_author(state, pubkey) {
400 tracing::debug!(
401 "Blossom write allowed for {}... (allowed writer)",
402 &pubkey[..8.min(pubkey.len())]
403 );
404 return Ok(());
405 }
406
407 tracing::info!(
409 "Blossom write denied for {}... (not in allowed_npubs or social graph)",
410 &pubkey[..8.min(pubkey.len())]
411 );
412 Err(Response::builder()
413 .status(StatusCode::FORBIDDEN)
414 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
415 .header(header::CONTENT_TYPE, "application/json")
416 .body(Body::from(
417 r#"{"error":"Write access denied. Your pubkey is not in the allowed list."}"#,
418 ))
419 .unwrap())
420}
421
422fn is_allowed_write_author(state: &AppState, pubkey: &str) -> bool {
423 if state.allowed_pubkeys.contains(pubkey) {
424 return true;
425 }
426
427 state
428 .social_graph
429 .as_ref()
430 .map(|sg| sg.check_write_access(pubkey))
431 .unwrap_or(false)
432}
433
434fn can_accept_upload_author(state: &AppState, pubkey: &str) -> bool {
435 state.public_writes || is_allowed_write_author(state, pubkey)
436}
437
438fn validate_upload_payload(
439 body: &[u8],
440 content_type: &str,
441 can_upload_author: bool,
442 require_random_untrusted_ingest: bool,
443) -> Result<(), (StatusCode, String)> {
444 let is_chk_upload = is_chk_content_type(content_type);
445
446 if !is_chk_upload && !can_upload_author {
447 return Err((
448 StatusCode::FORBIDDEN,
449 "Raw media uploads require write access".to_string(),
450 ));
451 }
452
453 if is_chk_upload {
454 let require_random = require_random_untrusted_ingest && !can_upload_author;
455 validate_untrusted_blob(body, require_random)
456 .map_err(|IngestRejection { status, reason }| (status, reason))?;
457 }
458
459 Ok(())
460}
461
462fn blossom_json_error(status: StatusCode, reason: impl Into<String>) -> Response<Body> {
463 let reason = reason.into();
464 Response::builder()
465 .status(status)
466 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
467 .header("X-Reason", reason.as_str())
468 .header(header::CONTENT_TYPE, "application/json")
469 .body(Body::from(format!(r#"{{"error":"{}"}}"#, reason)))
470 .unwrap()
471}
472
473fn blossom_retryable_json_error(
474 status: StatusCode,
475 reason: impl Into<String>,
476 retry_after_seconds: u64,
477) -> Response<Body> {
478 let reason = reason.into();
479 Response::builder()
480 .status(status)
481 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
482 .header("X-Reason", reason.as_str())
483 .header(header::RETRY_AFTER, retry_after_seconds.to_string())
484 .header(header::CONTENT_TYPE, "application/json")
485 .body(Body::from(format!(r#"{{"error":"{}"}}"#, reason)))
486 .unwrap()
487}
488
489#[derive(Debug)]
490enum BlobWriteError {
491 Busy(&'static str),
492 Storage(anyhow::Error),
493}
494
495impl std::fmt::Display for BlobWriteError {
496 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
497 match self {
498 Self::Busy(reason) => f.write_str(reason),
499 Self::Storage(error) => write!(f, "{error}"),
500 }
501 }
502}
503
504impl std::error::Error for BlobWriteError {}
505
506impl From<anyhow::Error> for BlobWriteError {
507 fn from(error: anyhow::Error) -> Self {
508 Self::Storage(error)
509 }
510}
511
512fn blob_io_write_error(error: BlobIoTaskError) -> BlobWriteError {
513 if error.is_busy() {
514 BlobWriteError::Busy(BLOB_WRITE_BUSY)
515 } else {
516 BlobWriteError::Storage(anyhow::anyhow!(error))
517 }
518}
519
520fn blob_write_error_response(error: BlobWriteError) -> Response<Body> {
521 match error {
522 BlobWriteError::Busy(reason) => {
523 blossom_retryable_json_error(StatusCode::SERVICE_UNAVAILABLE, reason, 2)
524 }
525 BlobWriteError::Storage(error) => Response::builder()
526 .status(StatusCode::INTERNAL_SERVER_ERROR)
527 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
528 .header("X-Reason", "Storage error")
529 .header(header::CONTENT_TYPE, "application/json")
530 .body(Body::from(format!(r#"{{"error":"{}"}}"#, error)))
531 .unwrap(),
532 }
533}
534
535#[derive(Debug, Clone, Serialize, Deserialize)]
537pub struct BlobDescriptor {
538 pub url: String,
539 pub sha256: String,
540 pub size: u64,
541 #[serde(rename = "type")]
542 pub mime_type: String,
543 pub uploaded: u64,
544}
545
546#[derive(Debug, Deserialize)]
547pub struct BatchUploadBlob {
548 pub sha256: String,
549 #[serde(default, alias = "contentType")]
550 pub content_type: Option<String>,
551 pub data: String,
552}
553
554#[derive(Debug, Deserialize)]
555pub struct BatchUploadRequest {
556 pub blobs: Vec<BatchUploadBlob>,
557}
558
559#[derive(Debug, Serialize, Deserialize)]
560pub struct BatchUploadResponse {
561 pub uploaded: usize,
562 pub blobs: Vec<BlobDescriptor>,
563}
564
565#[derive(Debug)]
566struct DecodedBatchUploadBlob {
567 sha256: String,
568 content_type: Option<String>,
569 data: Vec<u8>,
570}
571
572#[derive(Debug, Deserialize)]
573pub struct UploadCheckRequest {
574 pub hashes: Vec<String>,
575}
576
577#[derive(Debug, Serialize, Deserialize)]
578pub struct UploadCheckResponse {
579 pub count: usize,
580 pub present: String,
581}
582
583#[derive(Debug, Deserialize)]
585pub struct ListQuery {
586 pub since: Option<u64>,
587 pub until: Option<u64>,
588 pub limit: Option<usize>,
589 pub cursor: Option<String>,
590}
591
592#[derive(Debug)]
594pub struct BlossomAuth {
595 pub pubkey: String,
596 pub kind: u16,
597 pub created_at: u64,
598 pub expiration: Option<u64>,
599 pub action: Option<String>, pub blob_hashes: Vec<String>, pub batch_hashes: Vec<String>, pub server: Option<String>, }
604
605pub fn verify_blossom_auth(
608 headers: &HeaderMap,
609 required_action: &str,
610 required_hash: Option<&str>,
611) -> Result<BlossomAuth, (StatusCode, &'static str)> {
612 let auth_header = headers
613 .get(header::AUTHORIZATION)
614 .and_then(|v| v.to_str().ok())
615 .ok_or((StatusCode::UNAUTHORIZED, "Missing Authorization header"))?;
616
617 let nostr_event = auth_header.strip_prefix("Nostr ").ok_or((
618 StatusCode::UNAUTHORIZED,
619 "Invalid auth scheme, expected 'Nostr'",
620 ))?;
621
622 let engine = base64::engine::general_purpose::STANDARD;
624 let event_bytes = engine
625 .decode(nostr_event)
626 .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid base64 in auth header"))?;
627
628 let event_json: serde_json::Value = serde_json::from_slice(&event_bytes)
629 .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid JSON in auth event"))?;
630
631 let kind = event_json["kind"]
633 .as_u64()
634 .ok_or((StatusCode::BAD_REQUEST, "Missing kind in event"))?;
635
636 if kind != BLOSSOM_AUTH_KIND as u64 {
637 return Err((
638 StatusCode::BAD_REQUEST,
639 "Invalid event kind, expected 24242",
640 ));
641 }
642
643 let pubkey = event_json["pubkey"]
644 .as_str()
645 .ok_or((StatusCode::BAD_REQUEST, "Missing pubkey in event"))?
646 .to_string();
647
648 let created_at = event_json["created_at"]
649 .as_u64()
650 .ok_or((StatusCode::BAD_REQUEST, "Missing created_at in event"))?;
651
652 let sig = event_json["sig"]
653 .as_str()
654 .ok_or((StatusCode::BAD_REQUEST, "Missing signature in event"))?;
655
656 if !verify_nostr_signature(&event_json, &pubkey, sig) {
658 return Err((StatusCode::UNAUTHORIZED, "Invalid signature"));
659 }
660
661 let tags = event_json["tags"]
663 .as_array()
664 .ok_or((StatusCode::BAD_REQUEST, "Missing tags in event"))?;
665
666 let mut expiration: Option<u64> = None;
667 let mut action: Option<String> = None;
668 let mut blob_hashes: Vec<String> = Vec::new();
669 let mut batch_hashes: Vec<String> = Vec::new();
670 let mut server: Option<String> = None;
671
672 for tag in tags {
673 let tag_arr = tag.as_array();
674 if let Some(arr) = tag_arr {
675 if arr.len() >= 2 {
676 let tag_name = arr[0].as_str().unwrap_or("");
677 let tag_value = arr[1].as_str().unwrap_or("");
678
679 match tag_name {
680 "t" => action = Some(tag_value.to_string()),
681 "x" => blob_hashes.push(tag_value.to_lowercase()),
682 BATCH_UPLOAD_HASH_LIST_AUTH_TAG => batch_hashes.push(tag_value.to_lowercase()),
683 "expiration" => expiration = tag_value.parse().ok(),
684 "server" => server = Some(tag_value.to_string()),
685 _ => {}
686 }
687 }
688 }
689 }
690
691 let now = SystemTime::now()
693 .duration_since(UNIX_EPOCH)
694 .unwrap()
695 .as_secs();
696
697 if let Some(exp) = expiration {
698 if exp < now {
699 return Err((StatusCode::UNAUTHORIZED, "Authorization expired"));
700 }
701 }
702
703 if created_at > now + 60 {
705 return Err((StatusCode::BAD_REQUEST, "Event created_at is in the future"));
706 }
707
708 if let Some(ref act) = action {
710 if act != required_action {
711 return Err((StatusCode::FORBIDDEN, "Action mismatch"));
712 }
713 } else {
714 return Err((StatusCode::BAD_REQUEST, "Missing 't' tag for action"));
715 }
716
717 if let Some(hash) = required_hash {
719 if !blob_hashes.is_empty() && !blob_hashes.contains(&hash.to_lowercase()) {
720 return Err((StatusCode::FORBIDDEN, "Blob hash not authorized"));
721 }
722 }
723
724 Ok(BlossomAuth {
725 pubkey,
726 kind: kind as u16,
727 created_at,
728 expiration,
729 action,
730 blob_hashes,
731 batch_hashes,
732 server,
733 })
734}
735
736fn verify_nostr_signature(event: &serde_json::Value, pubkey: &str, sig: &str) -> bool {
738 use secp256k1::{schnorr::Signature, Message, Secp256k1, XOnlyPublicKey};
739
740 let content = event["content"].as_str().unwrap_or("");
742 let full_serialized = format!(
743 "[0,\"{}\",{},{},{},\"{}\"]",
744 pubkey,
745 event["created_at"],
746 event["kind"],
747 event["tags"],
748 escape_json_string(content),
749 );
750
751 let mut hasher = Sha256::new();
752 hasher.update(full_serialized.as_bytes());
753 let event_id = hasher.finalize();
754
755 let pubkey_bytes = match hex::decode(pubkey) {
757 Ok(b) => b,
758 Err(_) => return false,
759 };
760
761 let sig_bytes = match hex::decode(sig) {
762 Ok(b) => b,
763 Err(_) => return false,
764 };
765
766 let secp = Secp256k1::verification_only();
767
768 let xonly_pubkey = match XOnlyPublicKey::from_slice(&pubkey_bytes) {
769 Ok(pk) => pk,
770 Err(_) => return false,
771 };
772
773 let signature = match Signature::from_slice(&sig_bytes) {
774 Ok(s) => s,
775 Err(_) => return false,
776 };
777
778 let message = match Message::from_digest_slice(&event_id) {
779 Ok(m) => m,
780 Err(_) => return false,
781 };
782
783 secp.verify_schnorr(&signature, &message, &xonly_pubkey)
784 .is_ok()
785}
786
787fn escape_json_string(s: &str) -> String {
789 let mut result = String::new();
790 for c in s.chars() {
791 match c {
792 '"' => result.push_str("\\\""),
793 '\\' => result.push_str("\\\\"),
794 '\n' => result.push_str("\\n"),
795 '\r' => result.push_str("\\r"),
796 '\t' => result.push_str("\\t"),
797 c if c.is_control() => {
798 result.push_str(&format!("\\u{:04x}", c as u32));
799 }
800 c => result.push(c),
801 }
802 }
803 result
804}
805
806pub async fn cors_preflight(headers: HeaderMap) -> impl IntoResponse {
809 let allowed_headers = headers
811 .get(header::ACCESS_CONTROL_REQUEST_HEADERS)
812 .and_then(|v| v.to_str().ok())
813 .unwrap_or("Authorization, Content-Type, X-SHA-256, x-sha-256");
814
815 let full_allowed = format!(
817 "{}, Authorization, Content-Type, X-SHA-256, x-sha-256, Accept, Cache-Control",
818 allowed_headers
819 );
820
821 Response::builder()
822 .status(StatusCode::NO_CONTENT)
823 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
824 .header(
825 header::ACCESS_CONTROL_ALLOW_METHODS,
826 "GET, HEAD, POST, PUT, DELETE, OPTIONS",
827 )
828 .header(header::ACCESS_CONTROL_ALLOW_HEADERS, full_allowed)
829 .header(header::ACCESS_CONTROL_MAX_AGE, "86400")
830 .body(Body::empty())
831 .unwrap()
832}
833
834pub async fn head_upload(State(state): State<AppState>, headers: HeaderMap) -> impl IntoResponse {
836 let sha256_hex = match first_header_value(&headers, "x-sha-256") {
837 Some(hash) if is_valid_sha256(&hash) => hash.to_ascii_lowercase(),
838 Some(_) => {
839 return Response::builder()
840 .status(StatusCode::BAD_REQUEST)
841 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
842 .header("X-Reason", "Invalid X-SHA-256 header")
843 .body(Body::empty())
844 .unwrap();
845 }
846 None => {
847 return Response::builder()
848 .status(StatusCode::BAD_REQUEST)
849 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
850 .header("X-Reason", "Missing X-SHA-256 header")
851 .body(Body::empty())
852 .unwrap();
853 }
854 };
855
856 let content_length = match first_header_value(&headers, "x-content-length")
857 .or_else(|| first_header_value(&headers, header::CONTENT_LENGTH.as_str()))
858 .and_then(|value| value.parse::<usize>().ok())
859 {
860 Some(length) => length,
861 None => {
862 return Response::builder()
863 .status(StatusCode::BAD_REQUEST)
864 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
865 .header("X-Reason", "Missing or invalid X-Content-Length header")
866 .body(Body::empty())
867 .unwrap();
868 }
869 };
870
871 if content_length > state.max_upload_bytes {
872 return Response::builder()
873 .status(StatusCode::PAYLOAD_TOO_LARGE)
874 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
875 .header("X-Reason", "Upload exceeds maximum size")
876 .body(Body::empty())
877 .unwrap();
878 }
879
880 let content_type = first_header_value(&headers, "x-content-type")
881 .or_else(|| first_header_value(&headers, header::CONTENT_TYPE.as_str()))
882 .unwrap_or_else(|| "application/octet-stream".to_string());
883
884 if !is_chk_content_type(&content_type_base(&content_type)) {
885 let auth = verify_blossom_auth(&headers, "upload", Some(&sha256_hex));
886 let can_upload_raw = auth
887 .as_ref()
888 .map(|auth| can_accept_upload_author(&state, &auth.pubkey))
889 .unwrap_or(false);
890 if !can_upload_raw {
891 let status = if auth.is_err() {
892 StatusCode::UNAUTHORIZED
893 } else {
894 StatusCode::FORBIDDEN
895 };
896 return Response::builder()
897 .status(status)
898 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
899 .header("X-Reason", "Raw media uploads require write access")
900 .body(Body::empty())
901 .unwrap();
902 }
903 }
904
905 Response::builder()
906 .status(StatusCode::OK)
907 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
908 .body(Body::empty())
909 .unwrap()
910}
911
912fn encode_upload_check_bitset(bits: &[bool]) -> String {
913 let mut bytes = vec![0u8; bits.len().div_ceil(8)];
914 for (index, present) in bits.iter().enumerate() {
915 if *present {
916 bytes[index / 8] |= 1 << (index % 8);
917 }
918 }
919 base64::engine::general_purpose::STANDARD.encode(bytes)
920}
921
922pub async fn upload_check(
924 State(state): State<AppState>,
925 Json(payload): Json<UploadCheckRequest>,
926) -> impl IntoResponse {
927 if payload.hashes.len() > MAX_UPLOAD_CHECK_HASHES {
928 return blossom_json_error(
929 StatusCode::PAYLOAD_TOO_LARGE,
930 format!("Too many hashes; maximum is {}", MAX_UPLOAD_CHECK_HASHES),
931 );
932 }
933
934 let mut requested = Vec::with_capacity(payload.hashes.len());
935 let mut unique = Vec::new();
936 for hash in payload.hashes {
937 let hash = hash.trim().to_ascii_lowercase();
938 if !is_valid_sha256(&hash) {
939 return blossom_json_error(StatusCode::BAD_REQUEST, "Invalid SHA256 hash");
940 }
941 let bytes: [u8; 32] = match from_hex(&hash) {
942 Ok(bytes) => bytes,
943 Err(_) => return blossom_json_error(StatusCode::BAD_REQUEST, "Invalid SHA256 hash"),
944 };
945 requested.push(bytes);
946 unique.push(bytes);
947 }
948
949 unique.sort_unstable();
950 unique.dedup();
951
952 let existing = if unique.is_empty() {
953 Vec::new()
954 } else {
955 let store = state.store.clone();
956 let lookup_hashes = unique.clone();
957 match run_blob_metadata_read(move || {
958 store
959 .router()
960 .existing_local_hashes_in_sorted_candidates(&lookup_hashes)
961 })
962 .await
963 {
964 Ok(Ok(existing)) => existing,
965 Ok(Err(error)) => {
966 tracing::debug!("Blossom upload check failed: {}", error);
967 return blossom_json_error(StatusCode::INTERNAL_SERVER_ERROR, "Storage error");
968 }
969 Err(error) if error.is_busy() => {
970 return blossom_retryable_json_error(
971 StatusCode::SERVICE_UNAVAILABLE,
972 BLOB_READ_BUSY,
973 1,
974 );
975 }
976 Err(error) if error.is_timeout() => {
977 return blossom_retryable_json_error(
978 StatusCode::SERVICE_UNAVAILABLE,
979 "Blob check timed out",
980 1,
981 );
982 }
983 Err(error) => {
984 tracing::debug!("Blossom upload check task failed: {}", error);
985 return blossom_json_error(StatusCode::INTERNAL_SERVER_ERROR, "Storage error");
986 }
987 }
988 };
989
990 let present_unique: HashSet<[u8; 32]> = unique
991 .into_iter()
992 .zip(existing)
993 .filter_map(|(hash, present)| present.then_some(hash))
994 .collect();
995 let present_bits: Vec<bool> = requested
996 .iter()
997 .map(|hash| present_unique.contains(hash))
998 .collect();
999
1000 Response::builder()
1001 .status(StatusCode::OK)
1002 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
1003 .header(header::CONTENT_TYPE, "application/json")
1004 .body(Body::from(
1005 serde_json::to_string(&UploadCheckResponse {
1006 count: present_bits.len(),
1007 present: encode_upload_check_bitset(&present_bits),
1008 })
1009 .unwrap(),
1010 ))
1011 .unwrap()
1012}
1013
1014pub async fn head_blob(
1016 State(state): State<AppState>,
1017 Path(id): Path<String>,
1018 connect_info: axum::extract::ConnectInfo<std::net::SocketAddr>,
1019) -> impl IntoResponse {
1020 let is_localhost = connect_info.0.ip().is_loopback();
1021 let (hash_part, ext) = parse_hash_and_extension(&id);
1022
1023 if !is_valid_sha256(hash_part) {
1024 return Response::builder()
1025 .status(StatusCode::BAD_REQUEST)
1026 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
1027 .header("X-Reason", "Invalid SHA256 hash")
1028 .body(Body::empty())
1029 .unwrap();
1030 }
1031
1032 let sha256_hex = hash_part.to_lowercase();
1033 let sha256_bytes: [u8; 32] = match from_hex(&sha256_hex) {
1034 Ok(b) => b,
1035 Err(_) => {
1036 return Response::builder()
1037 .status(StatusCode::BAD_REQUEST)
1038 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
1039 .header("X-Reason", "Invalid SHA256 format")
1040 .body(Body::empty())
1041 .unwrap();
1042 }
1043 };
1044
1045 let blob_size = if let Some(cached) = state.blob_cache.get_size(&sha256_hex) {
1050 Ok(Ok(cached))
1051 } else {
1052 let store = state.store.clone();
1053 let result = match run_blob_metadata_read(move || store.blob_size(&sha256_bytes)).await {
1054 Ok(result) => Ok(result),
1055 Err(error) if error.is_busy() => {
1056 return Response::builder()
1057 .status(StatusCode::SERVICE_UNAVAILABLE)
1058 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
1059 .header(header::CACHE_CONTROL, NOT_FOUND_CACHE_CONTROL)
1060 .header("Retry-After", "1")
1061 .header("X-Reason", BLOB_READ_BUSY)
1062 .body(Body::empty())
1063 .unwrap();
1064 }
1065 Err(_) => Err(()),
1066 };
1067 if let Ok(Ok(size)) = result {
1068 state.blob_cache.put_size(sha256_hex.clone(), size);
1069 }
1070 result
1071 };
1072
1073 match blob_size {
1074 Ok(Ok(Some(size))) => {
1075 let mime_type = ext
1076 .map(|e| get_mime_type(&format!("file{}", e)))
1077 .unwrap_or("application/octet-stream");
1078
1079 let mut builder = Response::builder()
1080 .status(StatusCode::OK)
1081 .header(header::CONTENT_TYPE, mime_type)
1082 .header(header::CONTENT_LENGTH, size)
1083 .header(header::ACCEPT_RANGES, "bytes")
1084 .header(header::CACHE_CONTROL, IMMUTABLE_CACHE_CONTROL)
1085 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*");
1086 if is_localhost {
1087 builder = builder.header("X-Source", "local");
1088 }
1089 builder.body(Body::empty()).unwrap()
1090 }
1091 Ok(Ok(None)) => Response::builder()
1092 .status(StatusCode::NOT_FOUND)
1093 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
1094 .header(header::CACHE_CONTROL, IMMUTABLE_NOT_FOUND_CACHE_CONTROL)
1095 .header("X-Reason", "Blob not found")
1096 .body(Body::empty())
1097 .unwrap(),
1098 _ => Response::builder()
1099 .status(StatusCode::INTERNAL_SERVER_ERROR)
1100 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
1101 .body(Body::empty())
1102 .unwrap(),
1103 }
1104}
1105
1106async fn store_blossom_blob_without_blocking_runtime(
1107 state: &AppState,
1108 data: axum::body::Bytes,
1109 pubkey: [u8; 32],
1110 track_ownership: bool,
1111) -> Result<bool, BlobWriteError> {
1112 let mut hasher = Sha256::new();
1113 hasher.update(&data);
1114 let hash_hex = hex::encode(hasher.finalize());
1115 let data_for_cache = data.clone();
1116 let store = state.store.clone();
1117 let inserted = run_blob_write(move || {
1118 let inserted = if track_ownership {
1119 store.put_owned_blob_with_inserted(&data, &pubkey)?.1
1120 } else {
1121 store.put_cached_blob_with_inserted(&data)?.1
1122 };
1123 Ok::<_, anyhow::Error>(inserted)
1124 })
1125 .await
1126 .map_err(blob_io_write_error)??;
1127 state
1128 .blob_cache
1129 .put_size(hash_hex.clone(), Some(data_for_cache.len() as u64));
1130 state.blob_cache.put_body(hash_hex, &data_for_cache);
1131 Ok(inserted)
1132}
1133
1134fn upload_descriptor_response(status: StatusCode, descriptor: &BlobDescriptor) -> Response<Body> {
1135 Response::builder()
1136 .status(status)
1137 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
1138 .header(header::CONTENT_TYPE, "application/json")
1139 .body(Body::from(serde_json::to_string(descriptor).unwrap()))
1140 .unwrap()
1141}
1142
1143fn make_blob_descriptor(
1144 headers: &HeaderMap,
1145 sha256_hex: String,
1146 size: u64,
1147 mime_type: String,
1148 uploaded: u64,
1149) -> BlobDescriptor {
1150 let url = blossom_blob_url(headers, &sha256_hex, &mime_type);
1151 BlobDescriptor {
1152 url,
1153 sha256: sha256_hex,
1154 size,
1155 mime_type,
1156 uploaded,
1157 }
1158}
1159
1160fn blossom_blob_url(headers: &HeaderMap, sha256_hex: &str, mime_type: &str) -> String {
1161 format!(
1162 "{}/{}{}",
1163 blossom_public_base_url(headers),
1164 sha256_hex,
1165 descriptor_extension_for_mime(mime_type)
1166 )
1167}
1168
1169fn descriptor_extension_for_mime(mime_type: &str) -> &'static str {
1170 let base = content_type_base(mime_type);
1171 mime_to_extension(&base)
1172}
1173
1174fn blossom_public_base_url(headers: &HeaderMap) -> String {
1175 if let Some(configured) = configured_public_base_url() {
1176 return configured;
1177 }
1178
1179 let host = first_header_value(headers, "x-forwarded-host")
1180 .and_then(|value| normalize_host(&value))
1181 .or_else(|| {
1182 forwarded_header_param(headers, "host").and_then(|value| normalize_host(&value))
1183 })
1184 .or_else(|| {
1185 first_header_value(headers, header::HOST.as_str())
1186 .and_then(|value| normalize_host(&value))
1187 });
1188
1189 let Some(host) = host else {
1190 return "http://localhost".to_string();
1191 };
1192
1193 let scheme = first_header_value(headers, "x-forwarded-proto")
1194 .and_then(|value| normalize_scheme(&value))
1195 .or_else(|| {
1196 forwarded_header_param(headers, "proto").and_then(|value| normalize_scheme(&value))
1197 })
1198 .or_else(|| cloudflare_visitor_scheme(headers))
1199 .unwrap_or_else(|| default_scheme_for_host(&host).to_string());
1200
1201 format!("{scheme}://{host}")
1202}
1203
1204fn configured_public_base_url() -> Option<String> {
1205 [
1206 BLOSSOM_PUBLIC_BASE_URL_ENV,
1207 LEGACY_BLOSSOM_PUBLIC_BASE_URL_ENV,
1208 ]
1209 .into_iter()
1210 .find_map(|name| {
1211 std::env::var(name)
1212 .ok()
1213 .and_then(|value| normalize_public_base_url(&value))
1214 })
1215}
1216
1217fn normalize_public_base_url(value: &str) -> Option<String> {
1218 let trimmed = value.trim().trim_end_matches('/');
1219 if (trimmed.starts_with("https://") || trimmed.starts_with("http://"))
1220 && !trimmed.contains('?')
1221 && !trimmed.contains('#')
1222 {
1223 Some(trimmed.to_string())
1224 } else {
1225 None
1226 }
1227}
1228
1229fn first_header_value(headers: &HeaderMap, name: &str) -> Option<String> {
1230 let raw = headers.get(name)?.to_str().ok()?;
1231 first_header_component(raw)
1232}
1233
1234fn first_header_component(value: &str) -> Option<String> {
1235 clean_header_value(value.split(',').next().unwrap_or(value))
1236}
1237
1238fn forwarded_header_param(headers: &HeaderMap, name: &str) -> Option<String> {
1239 let raw = headers.get("forwarded")?.to_str().ok()?;
1240 let first = raw.split(',').next().unwrap_or(raw);
1241 for part in first.split(';') {
1242 let Some((key, value)) = part.split_once('=') else {
1243 continue;
1244 };
1245 if key.trim().eq_ignore_ascii_case(name) {
1246 return clean_header_value(value);
1247 }
1248 }
1249 None
1250}
1251
1252fn clean_header_value(value: &str) -> Option<String> {
1253 let trimmed = value.trim().trim_matches('"').trim();
1254 if trimmed.is_empty() || trimmed.chars().any(|ch| ch.is_control()) {
1255 None
1256 } else {
1257 Some(trimmed.to_string())
1258 }
1259}
1260
1261fn normalize_scheme(value: &str) -> Option<String> {
1262 match value.trim().trim_matches('"').to_ascii_lowercase().as_str() {
1263 "http" => Some("http".to_string()),
1264 "https" => Some("https".to_string()),
1265 _ => None,
1266 }
1267}
1268
1269fn normalize_host(value: &str) -> Option<String> {
1270 let mut host = value.trim().trim_matches('"').trim();
1271 if let Some(rest) = host.strip_prefix("http://") {
1272 host = rest;
1273 } else if let Some(rest) = host.strip_prefix("https://") {
1274 host = rest;
1275 }
1276 host = host.split('/').next().unwrap_or(host).trim();
1277 if host.is_empty()
1278 || host
1279 .chars()
1280 .any(|ch| ch.is_control() || ch.is_ascii_whitespace() || ch == '\\')
1281 {
1282 None
1283 } else {
1284 Some(host.to_string())
1285 }
1286}
1287
1288fn cloudflare_visitor_scheme(headers: &HeaderMap) -> Option<String> {
1289 let raw = headers.get("cf-visitor")?.to_str().ok()?;
1290 let parsed: serde_json::Value = serde_json::from_str(raw).ok()?;
1291 parsed
1292 .get("scheme")
1293 .and_then(|value| value.as_str())
1294 .and_then(normalize_scheme)
1295}
1296
1297fn default_scheme_for_host(host: &str) -> &'static str {
1298 let host_without_port = host
1299 .trim_start_matches('[')
1300 .split(']')
1301 .next()
1302 .unwrap_or(host)
1303 .split(':')
1304 .next()
1305 .unwrap_or(host)
1306 .to_ascii_lowercase();
1307 if host_without_port == "localhost"
1308 || host_without_port == "::1"
1309 || host_without_port.starts_with("127.")
1310 {
1311 "http"
1312 } else {
1313 "https"
1314 }
1315}
1316
1317fn optimistic_upload_queue_timeout() -> Duration {
1318 let millis = std::env::var(OPTIMISTIC_UPLOAD_QUEUE_TIMEOUT_MS_ENV)
1319 .ok()
1320 .and_then(|value| value.parse::<u64>().ok())
1321 .filter(|value| *value > 0)
1322 .unwrap_or(DEFAULT_OPTIMISTIC_UPLOAD_QUEUE_TIMEOUT_MS);
1323 Duration::from_millis(millis)
1324}
1325
1326fn optimistic_upload_inflight() -> &'static Mutex<HashSet<String>> {
1327 static INFLIGHT: OnceLock<Mutex<HashSet<String>>> = OnceLock::new();
1328 INFLIGHT.get_or_init(|| Mutex::new(HashSet::new()))
1329}
1330
1331fn optimistic_upload_is_inflight(hash_hex: &str) -> bool {
1332 optimistic_upload_inflight()
1333 .lock()
1334 .is_ok_and(|inflight| inflight.contains(hash_hex))
1335}
1336
1337fn mark_optimistic_upload_inflight(hash_hex: &str) -> bool {
1338 optimistic_upload_inflight()
1339 .lock()
1340 .map(|mut inflight| inflight.insert(hash_hex.to_string()))
1341 .unwrap_or(true)
1342}
1343
1344fn clear_optimistic_upload_inflight(hash_hex: &str) {
1345 if let Ok(mut inflight) = optimistic_upload_inflight().lock() {
1346 inflight.remove(hash_hex);
1347 }
1348}
1349
1350pub(super) fn optimistic_upload_queue_snapshot(state: &AppState) -> OptimisticUploadQueueSnapshot {
1351 let max_bytes = state.optimistic_upload_queue_bytes;
1352 let available_bytes = state
1353 .optimistic_upload_queue
1354 .available_permits()
1355 .min(max_bytes);
1356 let in_flight = optimistic_upload_inflight()
1357 .lock()
1358 .map(|inflight| inflight.len())
1359 .unwrap_or(0);
1360
1361 OptimisticUploadQueueSnapshot {
1362 enabled: state.optimistic_blossom_uploads,
1363 max_bytes,
1364 available_bytes,
1365 reserved_bytes: max_bytes.saturating_sub(available_bytes),
1366 in_flight,
1367 queue_timeout_ms: duration_millis_u64(optimistic_upload_queue_timeout()),
1368 }
1369}
1370
1371fn duration_millis_u64(duration: Duration) -> u64 {
1372 duration.as_millis().min(u128::from(u64::MAX)) as u64
1373}
1374
1375fn replica_item(hash: String, data: Vec<u8>, content_type: String) -> BatchUploadItem {
1376 BatchUploadItem {
1377 hash,
1378 data,
1379 content_type: Some(content_type),
1380 }
1381}
1382
1383fn prepare_blossom_upload_replication(
1384 state: &AppState,
1385 total_bytes: usize,
1386) -> Option<PreparedBlossomUploadReplication> {
1387 if state.blossom_upload_replicas.is_empty() {
1388 return None;
1389 }
1390 let permits = match u32::try_from(total_bytes.max(1)) {
1391 Ok(permits) => permits,
1392 Err(_) => {
1393 blossom_upload_replica_metrics()
1394 .skipped_jobs
1395 .fetch_add(1, Ordering::Relaxed);
1396 tracing::warn!(
1397 total_bytes,
1398 "Skipping Blossom write-behind replication because batch is too large"
1399 );
1400 return None;
1401 }
1402 };
1403 let permit = match state
1404 .blossom_upload_replica_queue
1405 .clone()
1406 .try_acquire_many_owned(permits)
1407 {
1408 Ok(permit) => permit,
1409 Err(error) => {
1410 blossom_upload_replica_metrics()
1411 .skipped_jobs
1412 .fetch_add(1, Ordering::Relaxed);
1413 tracing::warn!(
1414 total_bytes,
1415 targets = state.blossom_upload_replicas.len(),
1416 error = %error,
1417 "Skipping Blossom write-behind replication because queue is full"
1418 );
1419 return None;
1420 }
1421 };
1422
1423 let keys = match state.blossom_upload_replica_keys.clone() {
1424 Some(keys) => keys,
1425 None => {
1426 blossom_upload_replica_metrics()
1427 .skipped_jobs
1428 .fetch_add(1, Ordering::Relaxed);
1429 tracing::warn!(
1430 "Skipping Blossom write-behind replication because server keys are unavailable"
1431 );
1432 return None;
1433 }
1434 };
1435 Some(PreparedBlossomUploadReplication {
1436 servers: state.blossom_upload_replicas.clone(),
1437 keys,
1438 permit,
1439 total_bytes,
1440 })
1441}
1442
1443fn schedule_prepared_blossom_upload_replication(
1444 state: &AppState,
1445 prepared: PreparedBlossomUploadReplication,
1446 items: Vec<BatchUploadItem>,
1447) {
1448 if items.is_empty() {
1449 blossom_upload_replica_metrics()
1450 .skipped_jobs
1451 .fetch_add(1, Ordering::Relaxed);
1452 return;
1453 }
1454
1455 let job = BlossomReplicaUploadJob { prepared, items };
1456 let job = match state.blossom_upload_replica_scheduler.schedule(job) {
1457 Ok(()) => return,
1458 Err(job) => job,
1459 };
1460 spawn_blossom_replica_upload_batch(BlossomReplicaUploadBatch::from_job(job));
1461}
1462
1463async fn blossom_replica_coalescer_worker(mut receiver: mpsc::Receiver<BlossomReplicaUploadJob>) {
1464 let max_blobs = blossom_replica_coalesce_max_blobs();
1465 let max_bytes = blossom_replica_coalesce_max_bytes();
1466 let flush_delay = blossom_replica_coalesce_flush_delay();
1467
1468 while let Some(job) = receiver.recv().await {
1469 blossom_upload_replica_metrics()
1470 .coalesce_queued_jobs
1471 .fetch_sub(1, Ordering::Relaxed);
1472 let mut batch = BlossomReplicaUploadBatch::from_job(job);
1473 loop {
1474 if batch.reached_limits(max_blobs, max_bytes) {
1475 break;
1476 }
1477 match tokio::time::timeout(flush_delay, receiver.recv()).await {
1478 Ok(Some(next_job)) => {
1479 blossom_upload_replica_metrics()
1480 .coalesce_queued_jobs
1481 .fetch_sub(1, Ordering::Relaxed);
1482 if batch.can_append(&next_job, max_blobs, max_bytes) {
1483 batch.append(next_job);
1484 } else {
1485 spawn_blossom_replica_upload_batch(batch);
1486 batch = BlossomReplicaUploadBatch::from_job(next_job);
1487 }
1488 }
1489 Ok(None) => {
1490 spawn_blossom_replica_upload_batch(batch);
1491 return;
1492 }
1493 Err(_) => break,
1494 }
1495 }
1496 spawn_blossom_replica_upload_batch(batch);
1497 }
1498}
1499
1500fn spawn_blossom_replica_upload_batch(batch: BlossomReplicaUploadBatch) {
1501 tokio::spawn(async move {
1502 let _in_flight = BlossomReplicaInFlightGuard::new();
1503 let BlossomReplicaUploadBatch {
1504 servers,
1505 keys,
1506 permits: _permits,
1507 total_bytes,
1508 data_bytes: _,
1509 items,
1510 } = batch;
1511 let upload_limit = blossom_replica_upload_semaphore();
1512 let _upload_permit = match upload_limit.acquire().await {
1513 Ok(permit) => permit,
1514 Err(error) => {
1515 blossom_upload_replica_metrics()
1516 .failed_batches
1517 .fetch_add(1, Ordering::Relaxed);
1518 tracing::warn!(
1519 error = %error,
1520 "Skipping Blossom write-behind replication because upload limiter is closed"
1521 );
1522 return;
1523 }
1524 };
1525 let attempts = blossom_replica_upload_attempts();
1526 let client = BlossomClient::new((*keys).clone()).with_write_servers(servers.clone());
1527 for server in servers {
1528 for attempt in 1..=attempts {
1529 match client.upload_batch_to_server(&server, &items).await {
1530 Ok(Some(result)) => {
1531 let metrics = blossom_upload_replica_metrics();
1532 metrics.accepted_batches.fetch_add(1, Ordering::Relaxed);
1533 metrics
1534 .accepted_blobs
1535 .fetch_add(result.accepted as u64, Ordering::Relaxed);
1536 metrics
1537 .uploaded_blobs
1538 .fetch_add(result.uploaded as u64, Ordering::Relaxed);
1539 metrics
1540 .replicated_bytes
1541 .fetch_add(total_bytes as u64, Ordering::Relaxed);
1542 tracing::debug!(
1543 target = %server,
1544 accepted = result.accepted,
1545 uploaded = result.uploaded,
1546 total = items.len(),
1547 bytes = total_bytes,
1548 "Replicated Blossom upload batch"
1549 );
1550 break;
1551 }
1552 Ok(None) => {
1553 replicate_items_individually(&client, &server, &items, total_bytes).await;
1554 break;
1555 }
1556 Err(error) if attempt < attempts => {
1557 tracing::warn!(
1558 target = %server,
1559 error = %error,
1560 attempt,
1561 attempts,
1562 total = items.len(),
1563 bytes = total_bytes,
1564 "Blossom write-behind replication retrying"
1565 );
1566 tokio::time::sleep(Duration::from_millis(250 * attempt as u64)).await;
1567 }
1568 Err(error) => {
1569 blossom_upload_replica_metrics()
1570 .failed_batches
1571 .fetch_add(1, Ordering::Relaxed);
1572 tracing::warn!(
1573 target = %server,
1574 error = %error,
1575 attempt,
1576 attempts,
1577 total = items.len(),
1578 bytes = total_bytes,
1579 "Blossom write-behind replication failed"
1580 );
1581 break;
1582 }
1583 }
1584 }
1585 }
1586 });
1587}
1588
1589async fn replicate_items_individually(
1590 client: &BlossomClient,
1591 server: &str,
1592 items: &[BatchUploadItem],
1593 total_bytes: usize,
1594) {
1595 let mut uploaded = 0usize;
1596 let mut skipped = 0usize;
1597 let mut failed = 0usize;
1598 let server_list = [server.to_string()];
1599 for item in items {
1600 match client
1601 .upload_to_selected_servers(&item.data, &server_list)
1602 .await
1603 {
1604 Ok((_hash, successes)) if successes > 0 => uploaded += 1,
1605 Ok((_hash, _)) => skipped += 1,
1606 Err(error) => {
1607 failed += 1;
1608 tracing::warn!(
1609 target = %server,
1610 hash = %item.hash,
1611 error = %error,
1612 "Blossom write-behind item replication failed"
1613 );
1614 }
1615 }
1616 }
1617 let metrics = blossom_upload_replica_metrics();
1618 metrics.fallback_batches.fetch_add(1, Ordering::Relaxed);
1619 metrics
1620 .fallback_uploaded_blobs
1621 .fetch_add(uploaded as u64, Ordering::Relaxed);
1622 metrics
1623 .fallback_failed_blobs
1624 .fetch_add(failed as u64, Ordering::Relaxed);
1625 if failed == 0 {
1626 metrics.accepted_batches.fetch_add(1, Ordering::Relaxed);
1627 metrics
1628 .accepted_blobs
1629 .fetch_add(items.len() as u64, Ordering::Relaxed);
1630 metrics
1631 .uploaded_blobs
1632 .fetch_add(uploaded as u64, Ordering::Relaxed);
1633 metrics
1634 .replicated_bytes
1635 .fetch_add(total_bytes as u64, Ordering::Relaxed);
1636 } else {
1637 metrics.failed_batches.fetch_add(1, Ordering::Relaxed);
1638 }
1639 tracing::debug!(
1640 target = %server,
1641 uploaded,
1642 skipped,
1643 failed,
1644 total = items.len(),
1645 bytes = total_bytes,
1646 "Replicated Blossom upload items without batch support"
1647 );
1648}
1649
1650async fn acquire_optimistic_upload_queue(
1651 state: &AppState,
1652 permits: u32,
1653) -> Result<tokio::sync::OwnedSemaphorePermit, &'static str> {
1654 match tokio::time::timeout(
1655 optimistic_upload_queue_timeout(),
1656 state
1657 .optimistic_upload_queue
1658 .clone()
1659 .acquire_many_owned(permits),
1660 )
1661 .await
1662 {
1663 Ok(Ok(permit)) => Ok(permit),
1664 Ok(Err(_)) => Err("Optimistic upload queue is closed"),
1665 Err(_) => Err("Optimistic upload queue is full"),
1666 }
1667}
1668
1669pub async fn upload_blob(
1671 State(state): State<AppState>,
1672 headers: HeaderMap,
1673 body: axum::body::Bytes,
1674) -> impl IntoResponse {
1675 let max_size = state.max_upload_bytes;
1677 if body.len() > max_size {
1678 return Response::builder()
1679 .status(StatusCode::PAYLOAD_TOO_LARGE)
1680 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
1681 .header(header::CONTENT_TYPE, "application/json")
1682 .body(Body::from(format!(
1683 r#"{{"error":"Upload size {} bytes exceeds maximum {} bytes ({} MB)"}}"#,
1684 body.len(),
1685 max_size,
1686 max_size / 1024 / 1024
1687 )))
1688 .unwrap();
1689 }
1690
1691 let auth = match verify_blossom_auth(&headers, "upload", None) {
1693 Ok(a) => a,
1694 Err((status, reason)) => {
1695 return Response::builder()
1696 .status(status)
1697 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
1698 .header("X-Reason", reason)
1699 .header(header::CONTENT_TYPE, "application/json")
1700 .body(Body::from(format!(r#"{{"error":"{}"}}"#, reason)))
1701 .unwrap();
1702 }
1703 };
1704
1705 let content_type = headers
1707 .get(header::CONTENT_TYPE)
1708 .and_then(|v| v.to_str().ok())
1709 .unwrap_or("application/octet-stream")
1710 .to_string();
1711
1712 let is_allowed_author = is_allowed_write_author(&state, &auth.pubkey);
1714 let can_upload = can_accept_upload_author(&state, &auth.pubkey);
1715 if !can_upload {
1716 let _ = check_write_access(&state, &auth.pubkey);
1717 return Response::builder()
1718 .status(StatusCode::FORBIDDEN)
1719 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
1720 .header(header::CONTENT_TYPE, "application/json")
1721 .body(Body::from(r#"{"error":"Write access denied. Your pubkey is not in the allowed list and public writes are disabled."}"#))
1722 .unwrap();
1723 }
1724
1725 if let Err((status, reason)) = validate_upload_payload(
1726 &body,
1727 &content_type,
1728 can_upload,
1729 state.require_random_untrusted_ingest,
1730 ) {
1731 return blossom_json_error(status, reason);
1732 }
1733
1734 let mut hasher = Sha256::new();
1736 hasher.update(&body);
1737 let sha256_hash: [u8; 32] = hasher.finalize().into();
1738 let sha256_hex = hex::encode(sha256_hash);
1739
1740 if !auth.blob_hashes.is_empty() && !auth.blob_hashes.contains(&sha256_hex) {
1742 return Response::builder()
1743 .status(StatusCode::FORBIDDEN)
1744 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
1745 .header(
1746 "X-Reason",
1747 "Uploaded blob hash does not match authorized hash",
1748 )
1749 .header(header::CONTENT_TYPE, "application/json")
1750 .body(Body::from(r#"{"error":"Hash mismatch"}"#))
1751 .unwrap();
1752 }
1753
1754 let pubkey_bytes = match from_hex(&auth.pubkey) {
1756 Ok(b) => b,
1757 Err(_) => {
1758 return Response::builder()
1759 .status(StatusCode::BAD_REQUEST)
1760 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
1761 .header("X-Reason", "Invalid pubkey format")
1762 .body(Body::empty())
1763 .unwrap();
1764 }
1765 };
1766
1767 let size = body.len() as u64;
1768 let now = SystemTime::now()
1769 .duration_since(UNIX_EPOCH)
1770 .unwrap()
1771 .as_secs();
1772 let descriptor = make_blob_descriptor(&headers, sha256_hex.clone(), size, content_type, now);
1773
1774 if state.optimistic_blossom_uploads && optimistic_upload_is_inflight(&sha256_hex) {
1775 return upload_descriptor_response(StatusCode::ACCEPTED, &descriptor);
1776 }
1777
1778 if state.optimistic_blossom_uploads {
1781 let queued_bytes = body.len().max(OPTIMISTIC_UPLOAD_MIN_QUEUE_CHARGE_BYTES);
1782 if queued_bytes <= state.optimistic_upload_queue_bytes {
1783 let permits = queued_bytes as u32;
1784 let marked_inflight = mark_optimistic_upload_inflight(&sha256_hex);
1785 if !marked_inflight {
1786 return upload_descriptor_response(StatusCode::ACCEPTED, &descriptor);
1787 }
1788 let permit = match acquire_optimistic_upload_queue(&state, permits).await {
1789 Ok(permit) => permit,
1790 Err(_) => {
1791 clear_optimistic_upload_inflight(&sha256_hex);
1792 return blossom_retryable_json_error(
1793 StatusCode::SERVICE_UNAVAILABLE,
1794 "Optimistic upload queue is full",
1795 2,
1796 );
1797 }
1798 };
1799 let state_for_write = state.clone();
1800 let hash_for_log = sha256_hex.clone();
1801 let replica_body = body.clone();
1802 let replica_content_type = descriptor.mime_type.clone();
1803 tokio::spawn(async move {
1804 let _permit = permit;
1805 match store_blossom_blob_without_blocking_runtime(
1806 &state_for_write,
1807 body,
1808 pubkey_bytes,
1809 is_allowed_author,
1810 )
1811 .await
1812 {
1813 Ok(inserted) => {
1814 if inserted {
1815 if let Some(replication) = prepare_blossom_upload_replication(
1816 &state_for_write,
1817 replica_body.len(),
1818 ) {
1819 let replication_item = replica_item(
1820 hash_for_log.clone(),
1821 replica_body.to_vec(),
1822 replica_content_type,
1823 );
1824 schedule_prepared_blossom_upload_replication(
1825 &state_for_write,
1826 replication,
1827 vec![replication_item],
1828 );
1829 }
1830 }
1831 }
1832 Err(error) => {
1833 tracing::error!(
1834 "Background Blossom storage failed for {}: {:#}",
1835 hash_for_log,
1836 error
1837 );
1838 }
1839 }
1840 clear_optimistic_upload_inflight(&hash_for_log);
1841 });
1842
1843 return upload_descriptor_response(StatusCode::ACCEPTED, &descriptor);
1844 }
1845
1846 tracing::warn!(
1847 "Blossom upload {} is larger than optimistic queue budget {}; storing synchronously",
1848 queued_bytes,
1849 state.optimistic_upload_queue_bytes
1850 );
1851 }
1852
1853 let store_result = store_blossom_blob_without_blocking_runtime(
1854 &state,
1855 body.clone(),
1856 pubkey_bytes,
1857 is_allowed_author,
1858 )
1859 .await;
1860
1861 match store_result {
1862 Ok(inserted) => {
1863 if inserted {
1864 if let Some(replication) = prepare_blossom_upload_replication(&state, body.len()) {
1865 let replication_item = replica_item(
1866 sha256_hex.clone(),
1867 body.to_vec(),
1868 descriptor.mime_type.clone(),
1869 );
1870 schedule_prepared_blossom_upload_replication(
1871 &state,
1872 replication,
1873 vec![replication_item],
1874 );
1875 }
1876 }
1877 upload_descriptor_response(
1878 if inserted {
1879 StatusCode::CREATED
1880 } else {
1881 StatusCode::OK
1882 },
1883 &descriptor,
1884 )
1885 }
1886 Err(error) => blob_write_error_response(error),
1887 }
1888}
1889
1890fn take_binary_batch_bytes<'a>(
1891 body: &'a [u8],
1892 cursor: &mut usize,
1893 len: usize,
1894) -> Result<&'a [u8], (StatusCode, String)> {
1895 let end = cursor.checked_add(len).ok_or_else(|| {
1896 (
1897 StatusCode::PAYLOAD_TOO_LARGE,
1898 "Binary batch field length overflow".to_string(),
1899 )
1900 })?;
1901 if end > body.len() {
1902 return Err((
1903 StatusCode::BAD_REQUEST,
1904 "Binary batch body is truncated".to_string(),
1905 ));
1906 }
1907 let slice = &body[*cursor..end];
1908 *cursor = end;
1909 Ok(slice)
1910}
1911
1912fn read_binary_batch_u16(body: &[u8], cursor: &mut usize) -> Result<u16, (StatusCode, String)> {
1913 let bytes = take_binary_batch_bytes(body, cursor, 2)?;
1914 Ok(u16::from_be_bytes([bytes[0], bytes[1]]))
1915}
1916
1917fn read_binary_batch_u32(body: &[u8], cursor: &mut usize) -> Result<u32, (StatusCode, String)> {
1918 let bytes = take_binary_batch_bytes(body, cursor, 4)?;
1919 Ok(u32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]))
1920}
1921
1922fn read_binary_batch_u64(body: &[u8], cursor: &mut usize) -> Result<u64, (StatusCode, String)> {
1923 let bytes = take_binary_batch_bytes(body, cursor, 8)?;
1924 Ok(u64::from_be_bytes([
1925 bytes[0], bytes[1], bytes[2], bytes[3], bytes[4], bytes[5], bytes[6], bytes[7],
1926 ]))
1927}
1928
1929fn parse_binary_batch_upload(
1930 body: &[u8],
1931) -> Result<Vec<DecodedBatchUploadBlob>, (StatusCode, String)> {
1932 let mut cursor = 0usize;
1933 let magic = take_binary_batch_bytes(body, &mut cursor, BINARY_BATCH_UPLOAD_MAGIC.len())?;
1934 if magic != BINARY_BATCH_UPLOAD_MAGIC {
1935 return Err((
1936 StatusCode::BAD_REQUEST,
1937 "Invalid binary batch magic".to_string(),
1938 ));
1939 }
1940
1941 let count = read_binary_batch_u32(body, &mut cursor)? as usize;
1942 if count == 0 {
1943 return Err((StatusCode::BAD_REQUEST, "Batch is empty".to_string()));
1944 }
1945 if count > MAX_BATCH_UPLOAD_BLOBS {
1946 return Err((
1947 StatusCode::PAYLOAD_TOO_LARGE,
1948 "Batch contains too many blobs".to_string(),
1949 ));
1950 }
1951
1952 let mut total_bytes = 0usize;
1953 let mut blobs = Vec::with_capacity(count);
1954 for _ in 0..count {
1955 let hash = take_binary_batch_bytes(body, &mut cursor, 32)?;
1956 let content_type_len = read_binary_batch_u16(body, &mut cursor)? as usize;
1957 if content_type_len > MAX_BINARY_BATCH_CONTENT_TYPE_BYTES {
1958 return Err((
1959 StatusCode::PAYLOAD_TOO_LARGE,
1960 "Binary batch content type is too long".to_string(),
1961 ));
1962 }
1963 let data_len =
1964 usize::try_from(read_binary_batch_u64(body, &mut cursor)?).map_err(|_| {
1965 (
1966 StatusCode::PAYLOAD_TOO_LARGE,
1967 "Binary batch blob is too large".to_string(),
1968 )
1969 })?;
1970 total_bytes = total_bytes.checked_add(data_len).ok_or_else(|| {
1971 (
1972 StatusCode::PAYLOAD_TOO_LARGE,
1973 "Batch exceeds maximum upload size".to_string(),
1974 )
1975 })?;
1976 if total_bytes > MAX_BATCH_UPLOAD_BYTES {
1977 return Err((
1978 StatusCode::PAYLOAD_TOO_LARGE,
1979 "Batch exceeds maximum upload size".to_string(),
1980 ));
1981 }
1982
1983 let content_type = if content_type_len == 0 {
1984 None
1985 } else {
1986 let bytes = take_binary_batch_bytes(body, &mut cursor, content_type_len)?;
1987 Some(
1988 std::str::from_utf8(bytes)
1989 .map_err(|_| {
1990 (
1991 StatusCode::BAD_REQUEST,
1992 "Binary batch content type is not UTF-8".to_string(),
1993 )
1994 })?
1995 .to_string(),
1996 )
1997 };
1998 let data = take_binary_batch_bytes(body, &mut cursor, data_len)?.to_vec();
1999 blobs.push(DecodedBatchUploadBlob {
2000 sha256: hex::encode(hash),
2001 content_type,
2002 data,
2003 });
2004 }
2005
2006 if cursor != body.len() {
2007 return Err((
2008 StatusCode::BAD_REQUEST,
2009 "Binary batch body has trailing bytes".to_string(),
2010 ));
2011 }
2012
2013 Ok(blobs)
2014}
2015
2016fn blossom_auth_error_response(status: StatusCode, reason: &'static str) -> Response<Body> {
2017 Response::builder()
2018 .status(status)
2019 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2020 .header("X-Reason", reason)
2021 .header(header::CONTENT_TYPE, "application/json")
2022 .body(Body::from(format!(r#"{{"error":"{}"}}"#, reason)))
2023 .unwrap()
2024}
2025
2026fn verify_upload_batch_auth(headers: &HeaderMap) -> Result<BlossomAuth, Box<Response<Body>>> {
2027 verify_blossom_auth(headers, "upload", None)
2028 .map_err(|(status, reason)| Box::new(blossom_auth_error_response(status, reason)))
2029}
2030
2031pub async fn require_upload_auth_middleware(request: Request, next: Next) -> Response<Body> {
2032 if let Err(response) = verify_upload_batch_auth(request.headers()) {
2033 return *response;
2034 }
2035 next.run(request).await
2036}
2037
2038async fn upload_decoded_blob_batch(
2039 state: AppState,
2040 headers: HeaderMap,
2041 auth: BlossomAuth,
2042 blobs: Vec<DecodedBatchUploadBlob>,
2043 started_at: Instant,
2044 encoding: &'static str,
2045) -> Response<Body> {
2046 let slow_log_ms = slow_batch_upload_log_ms();
2047 let payload_blobs = blobs.len();
2048 if blobs.is_empty() {
2049 return blossom_json_error(StatusCode::BAD_REQUEST, "Batch is empty");
2050 }
2051 if blobs.len() > MAX_BATCH_UPLOAD_BLOBS {
2052 return blossom_json_error(
2053 StatusCode::PAYLOAD_TOO_LARGE,
2054 "Batch contains too many blobs",
2055 );
2056 }
2057
2058 let auth_ms = started_at.elapsed().as_millis();
2059
2060 let is_allowed_author = is_allowed_write_author(&state, &auth.pubkey);
2061 let can_upload = can_accept_upload_author(&state, &auth.pubkey);
2062 if !can_upload {
2063 let _ = check_write_access(&state, &auth.pubkey);
2064 return blossom_json_error(StatusCode::FORBIDDEN, "Write access denied");
2065 }
2066
2067 let pubkey_bytes = match from_hex(&auth.pubkey) {
2068 Ok(bytes) => bytes,
2069 Err(_) => return blossom_json_error(StatusCode::BAD_REQUEST, "Invalid pubkey format"),
2070 };
2071
2072 let now = SystemTime::now()
2073 .duration_since(UNIX_EPOCH)
2074 .unwrap()
2075 .as_secs();
2076 let mut total_bytes = 0usize;
2077 let mut items = Vec::with_capacity(payload_blobs);
2078 let mut replica_specs = Vec::with_capacity(payload_blobs);
2079 let mut descriptors = Vec::with_capacity(payload_blobs);
2080 let mut decode_hash_ms = 0u128;
2081 let mut validate_ms = 0u128;
2082
2083 if auth.blob_hashes.is_empty() && !auth.batch_hashes.is_empty() {
2084 let batch_hashes = blobs
2085 .iter()
2086 .map(|blob| blob.sha256.to_lowercase())
2087 .collect::<Vec<_>>();
2088 let batch_digest =
2089 match batch_upload_hash_list_digest(batch_hashes.iter().map(String::as_str)) {
2090 Ok(digest) => digest,
2091 Err(_) => return blossom_json_error(StatusCode::BAD_REQUEST, "Invalid blob hash"),
2092 };
2093 if !auth.batch_hashes.contains(&batch_digest) {
2094 return blossom_json_error(
2095 StatusCode::FORBIDDEN,
2096 "Batch hash list does not match authorization",
2097 );
2098 }
2099 }
2100
2101 for blob in blobs {
2102 let decode_started = Instant::now();
2103 let sha256_hex = blob.sha256.to_lowercase();
2104 let expected_hash: [u8; 32] = match from_hex(&sha256_hex) {
2105 Ok(hash) => hash,
2106 Err(_) => return blossom_json_error(StatusCode::BAD_REQUEST, "Invalid blob hash"),
2107 };
2108 if !auth.blob_hashes.is_empty() && !auth.blob_hashes.contains(&sha256_hex) {
2109 return blossom_json_error(
2110 StatusCode::FORBIDDEN,
2111 "Uploaded blob hash does not match authorized hash",
2112 );
2113 }
2114
2115 let data = blob.data;
2116 if data.len() > state.max_upload_bytes {
2117 return blossom_json_error(
2118 StatusCode::PAYLOAD_TOO_LARGE,
2119 "Blob exceeds maximum upload size",
2120 );
2121 }
2122 total_bytes = total_bytes.saturating_add(data.len());
2123 if total_bytes > MAX_BATCH_UPLOAD_BYTES {
2124 return blossom_json_error(
2125 StatusCode::PAYLOAD_TOO_LARGE,
2126 "Batch exceeds maximum upload size",
2127 );
2128 }
2129
2130 let mut hasher = Sha256::new();
2131 hasher.update(&data);
2132 let actual_hash: [u8; 32] = hasher.finalize().into();
2133 if actual_hash != expected_hash {
2134 return blossom_json_error(StatusCode::FORBIDDEN, "Hash mismatch");
2135 }
2136 decode_hash_ms += decode_started.elapsed().as_millis();
2137
2138 let validate_started = Instant::now();
2139 let content_type = blob
2140 .content_type
2141 .as_deref()
2142 .map(content_type_base)
2143 .unwrap_or_else(|| "application/octet-stream".to_string());
2144 if let Err((status, reason)) = validate_upload_payload(
2145 &data,
2146 &content_type,
2147 can_upload,
2148 state.require_random_untrusted_ingest,
2149 ) {
2150 return blossom_json_error(status, reason);
2151 }
2152 validate_ms += validate_started.elapsed().as_millis();
2153
2154 descriptors.push(make_blob_descriptor(
2155 &headers,
2156 sha256_hex.clone(),
2157 data.len() as u64,
2158 content_type.clone(),
2159 now,
2160 ));
2161 replica_specs.push((sha256_hex, content_type));
2162 items.push((actual_hash, data));
2163 }
2164 let prepare_ms = started_at.elapsed().as_millis();
2165
2166 let store = state.store.clone();
2167 let store_started = Instant::now();
2168 let stored = run_blob_write(move || {
2169 let report = if is_allowed_author {
2170 store.put_owned_blobs_report(&items, &pubkey_bytes)
2171 } else {
2172 store.put_cached_blobs_report(&items)
2173 }?;
2174 Ok::<_, anyhow::Error>((report, items))
2175 })
2176 .await
2177 .map_err(blob_io_write_error);
2178 let store_ms = store_started.elapsed().as_millis();
2179 let total_ms = started_at.elapsed().as_millis();
2180
2181 match stored {
2182 Ok(Ok((report, items))) => {
2183 for ((sha256_hex, _), (_, data)) in replica_specs.iter().zip(&items) {
2184 state
2185 .blob_cache
2186 .put_size(sha256_hex.clone(), Some(data.len() as u64));
2187 state.blob_cache.put_body(sha256_hex.clone(), data);
2188 }
2189 let uploaded = report.inserted;
2190 if uploaded > 0 {
2191 let inserted: HashSet<_> = report.inserted_hashes.iter().copied().collect();
2192 let replica_items = replica_specs
2193 .iter()
2194 .zip(items.iter())
2195 .filter(|&((_, _), (hash, _))| inserted.contains(hash))
2196 .map(|((sha256_hex, content_type), (_, data))| {
2197 replica_item(sha256_hex.clone(), data.clone(), content_type.clone())
2198 })
2199 .collect::<Vec<_>>();
2200 let replication =
2201 usize::try_from(report.inserted_bytes)
2202 .ok()
2203 .and_then(|inserted_bytes| {
2204 prepare_blossom_upload_replication(&state, inserted_bytes)
2205 });
2206 if let Some(replication) = replication {
2207 schedule_prepared_blossom_upload_replication(
2208 &state,
2209 replication,
2210 replica_items,
2211 );
2212 }
2213 }
2214 if slow_log_ms.is_some_and(|threshold| total_ms >= threshold) {
2215 tracing::warn!(
2216 blobs = payload_blobs,
2217 uploaded,
2218 total_bytes,
2219 total_ms,
2220 auth_ms,
2221 prepare_ms,
2222 decode_hash_ms,
2223 validate_ms,
2224 store_ms,
2225 encoding,
2226 allowed_author = is_allowed_author,
2227 "slow Blossom batch upload"
2228 );
2229 }
2230 Response::builder()
2231 .status(StatusCode::OK)
2232 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2233 .header(header::CONTENT_TYPE, "application/json")
2234 .body(Body::from(
2235 serde_json::to_string(&BatchUploadResponse {
2236 uploaded,
2237 blobs: descriptors,
2238 })
2239 .unwrap(),
2240 ))
2241 .unwrap()
2242 }
2243 Ok(Err(error)) => blob_write_error_response(error.into()),
2244 Err(error) => blob_write_error_response(error),
2245 }
2246}
2247
2248pub async fn upload_blob_batch(
2250 State(state): State<AppState>,
2251 headers: HeaderMap,
2252 Json(payload): Json<BatchUploadRequest>,
2253) -> impl IntoResponse {
2254 let started_at = Instant::now();
2255 let auth = match verify_upload_batch_auth(&headers) {
2256 Ok(auth) => auth,
2257 Err(response) => return *response,
2258 };
2259 let mut blobs = Vec::with_capacity(payload.blobs.len());
2260 for blob in payload.blobs {
2261 let data = match base64::engine::general_purpose::STANDARD.decode(blob.data.as_bytes()) {
2262 Ok(data) => data,
2263 Err(_) => return blossom_json_error(StatusCode::BAD_REQUEST, "Invalid blob data"),
2264 };
2265 blobs.push(DecodedBatchUploadBlob {
2266 sha256: blob.sha256,
2267 content_type: blob.content_type,
2268 data,
2269 });
2270 }
2271 upload_decoded_blob_batch(state, headers, auth, blobs, started_at, "json").await
2272}
2273
2274pub async fn upload_blob_batch_binary(
2276 State(state): State<AppState>,
2277 headers: HeaderMap,
2278 body: Bytes,
2279) -> impl IntoResponse {
2280 let started_at = Instant::now();
2281 let auth = match verify_upload_batch_auth(&headers) {
2282 Ok(auth) => auth,
2283 Err(response) => return *response,
2284 };
2285 let blobs = match parse_binary_batch_upload(&body) {
2286 Ok(blobs) => blobs,
2287 Err((status, reason)) => return blossom_json_error(status, reason),
2288 };
2289 upload_decoded_blob_batch(state, headers, auth, blobs, started_at, "binary").await
2290}
2291
2292pub async fn delete_blob(
2295 State(state): State<AppState>,
2296 Path(id): Path<String>,
2297 headers: HeaderMap,
2298) -> impl IntoResponse {
2299 let (hash_part, _) = parse_hash_and_extension(&id);
2300
2301 if !is_valid_sha256(hash_part) {
2302 return Response::builder()
2303 .status(StatusCode::BAD_REQUEST)
2304 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2305 .header("X-Reason", "Invalid SHA256 hash")
2306 .body(Body::empty())
2307 .unwrap();
2308 }
2309
2310 let sha256_hex = hash_part.to_lowercase();
2311
2312 let sha256_bytes = match from_hex(&sha256_hex) {
2314 Ok(b) => b,
2315 Err(_) => {
2316 return Response::builder()
2317 .status(StatusCode::BAD_REQUEST)
2318 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2319 .header("X-Reason", "Invalid SHA256 hash format")
2320 .body(Body::empty())
2321 .unwrap();
2322 }
2323 };
2324
2325 let auth = match verify_blossom_auth(&headers, "delete", Some(&sha256_hex)) {
2327 Ok(a) => a,
2328 Err((status, reason)) => {
2329 return Response::builder()
2330 .status(status)
2331 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2332 .header("X-Reason", reason)
2333 .body(Body::empty())
2334 .unwrap();
2335 }
2336 };
2337
2338 let pubkey_bytes = match from_hex(&auth.pubkey) {
2340 Ok(b) => b,
2341 Err(_) => {
2342 return Response::builder()
2343 .status(StatusCode::BAD_REQUEST)
2344 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2345 .header("X-Reason", "Invalid pubkey format")
2346 .body(Body::empty())
2347 .unwrap();
2348 }
2349 };
2350
2351 let store = state.store.clone();
2354 let ownership = run_blob_metadata_read(move || {
2355 let is_owner = store.is_blob_owner(&sha256_bytes, &pubkey_bytes)?;
2356 let has_owners = if is_owner {
2357 true
2358 } else {
2359 store.blob_has_owners(&sha256_bytes)?
2360 };
2361 Ok::<_, anyhow::Error>((is_owner, has_owners))
2362 })
2363 .await;
2364 match ownership {
2365 Ok(Ok((true, _))) => {}
2366 Ok(Ok((false, true))) => {
2367 return Response::builder()
2368 .status(StatusCode::FORBIDDEN)
2369 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2370 .header("X-Reason", "Not a blob owner")
2371 .body(Body::empty())
2372 .unwrap();
2373 }
2374 Ok(Ok((false, false))) => {
2375 return Response::builder()
2376 .status(StatusCode::NOT_FOUND)
2377 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2378 .header(header::CACHE_CONTROL, NOT_FOUND_CACHE_CONTROL)
2379 .header("X-Reason", "Blob not found")
2380 .body(Body::empty())
2381 .unwrap();
2382 }
2383 Ok(Err(_)) | Err(_) => {
2384 return Response::builder()
2385 .status(StatusCode::INTERNAL_SERVER_ERROR)
2386 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2387 .body(Body::empty())
2388 .unwrap();
2389 }
2390 }
2391
2392 let store = state.store.clone();
2394 match run_blob_write(move || store.delete_blossom_blob(&sha256_bytes, &pubkey_bytes)).await {
2395 Ok(Ok(fully_deleted)) => {
2396 Response::builder()
2399 .status(StatusCode::OK)
2400 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2401 .header(
2402 "X-Blob-Deleted",
2403 if fully_deleted { "true" } else { "false" },
2404 )
2405 .body(Body::empty())
2406 .unwrap()
2407 }
2408 Ok(Err(error)) => Response::builder()
2409 .status(
2410 if error
2411 .to_string()
2412 .contains(crate::storage::POOL_MIGRATION_DELETE_DISABLED)
2413 {
2414 StatusCode::SERVICE_UNAVAILABLE
2415 } else {
2416 StatusCode::INTERNAL_SERVER_ERROR
2417 },
2418 )
2419 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2420 .body(Body::empty())
2421 .unwrap(),
2422 Err(_) => Response::builder()
2423 .status(StatusCode::INTERNAL_SERVER_ERROR)
2424 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2425 .body(Body::empty())
2426 .unwrap(),
2427 }
2428}
2429
2430pub async fn list_blobs(
2432 State(state): State<AppState>,
2433 Path(pubkey): Path<String>,
2434 Query(query): Query<ListQuery>,
2435 headers: HeaderMap,
2436) -> impl IntoResponse {
2437 if pubkey.len() != 64 || !pubkey.chars().all(|c| c.is_ascii_hexdigit()) {
2439 return Response::builder()
2440 .status(StatusCode::BAD_REQUEST)
2441 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2442 .header("X-Reason", "Invalid pubkey format")
2443 .header(header::CONTENT_TYPE, "application/json")
2444 .body(Body::from("[]"))
2445 .unwrap();
2446 }
2447
2448 let pubkey_hex = pubkey.to_lowercase();
2449 let pubkey_bytes: [u8; 32] = match from_hex(&pubkey_hex) {
2450 Ok(b) => b,
2451 Err(_) => {
2452 return Response::builder()
2453 .status(StatusCode::BAD_REQUEST)
2454 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2455 .header("X-Reason", "Invalid pubkey format")
2456 .header(header::CONTENT_TYPE, "application/json")
2457 .body(Body::from("[]"))
2458 .unwrap();
2459 }
2460 };
2461
2462 let auth = match verify_blossom_auth(&headers, "list", None) {
2463 Ok(auth) => auth,
2464 Err((status, reason)) => {
2465 return Response::builder()
2466 .status(status)
2467 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2468 .header("X-Reason", reason)
2469 .header(header::CONTENT_TYPE, "application/json")
2470 .body(Body::from("[]"))
2471 .unwrap();
2472 }
2473 };
2474
2475 if !auth.pubkey.eq_ignore_ascii_case(&pubkey_hex) {
2476 return Response::builder()
2477 .status(StatusCode::FORBIDDEN)
2478 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2479 .header("X-Reason", "Pubkey mismatch")
2480 .header(header::CONTENT_TYPE, "application/json")
2481 .body(Body::from("[]"))
2482 .unwrap();
2483 }
2484
2485 let store = state.store.clone();
2487 match run_blob_metadata_read(move || store.list_blobs_by_pubkey(&pubkey_bytes)).await {
2488 Ok(Ok(blobs)) => {
2489 let mut filtered: Vec<_> = blobs
2491 .into_iter()
2492 .filter(|b| {
2493 if let Some(since) = query.since {
2494 if b.uploaded < since {
2495 return false;
2496 }
2497 }
2498 if let Some(until) = query.until {
2499 if b.uploaded > until {
2500 return false;
2501 }
2502 }
2503 true
2504 })
2505 .collect();
2506
2507 filtered.sort_by_key(|descriptor| std::cmp::Reverse(descriptor.uploaded));
2509
2510 let limit = query.limit.unwrap_or(100).min(1000);
2512 filtered.truncate(limit);
2513
2514 let descriptors: Vec<_> = filtered
2515 .into_iter()
2516 .map(|mut descriptor| {
2517 descriptor.url =
2518 blossom_blob_url(&headers, &descriptor.sha256, &descriptor.mime_type);
2519 descriptor
2520 })
2521 .collect();
2522
2523 Response::builder()
2524 .status(StatusCode::OK)
2525 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2526 .header(header::CONTENT_TYPE, "application/json")
2527 .body(Body::from(serde_json::to_string(&descriptors).unwrap()))
2528 .unwrap()
2529 }
2530 Ok(Err(_)) | Err(_) => Response::builder()
2531 .status(StatusCode::INTERNAL_SERVER_ERROR)
2532 .header(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*")
2533 .header(header::CONTENT_TYPE, "application/json")
2534 .body(Body::from("[]"))
2535 .unwrap(),
2536 }
2537}
2538
2539fn parse_hash_and_extension(id: &str) -> (&str, Option<&str>) {
2542 if let Some(dot_pos) = id.rfind('.') {
2543 (&id[..dot_pos], Some(&id[dot_pos..]))
2544 } else {
2545 (id, None)
2546 }
2547}
2548
2549fn is_valid_sha256(s: &str) -> bool {
2550 s.len() == 64 && s.chars().all(|c| c.is_ascii_hexdigit())
2551}
2552
2553#[cfg(test)]
2554fn store_blossom_blob(
2555 state: &AppState,
2556 data: &[u8],
2557 _sha256: &[u8; 32],
2558 pubkey: &[u8; 32],
2559 track_ownership: bool,
2560) -> anyhow::Result<()> {
2561 if track_ownership {
2562 state.store.put_owned_blob(data, pubkey)?;
2563 } else {
2564 state.store.put_cached_blob(data)?;
2565 }
2566
2567 Ok(())
2568}
2569
2570fn mime_to_extension(mime: &str) -> &'static str {
2571 match mime {
2572 "image/png" => ".png",
2573 "image/jpeg" => ".jpg",
2574 "image/gif" => ".gif",
2575 "image/webp" => ".webp",
2576 "image/svg+xml" => ".svg",
2577 "video/mp4" => ".mp4",
2578 "video/webm" => ".webm",
2579 "audio/mpeg" => ".mp3",
2580 "audio/ogg" => ".ogg",
2581 "application/pdf" => ".pdf",
2582 "text/plain" => ".txt",
2583 "text/html" => ".html",
2584 "application/json" => ".json",
2585 _ => ".bin",
2586 }
2587}
2588
2589#[cfg(test)]
2590mod tests {
2591 use super::*;
2592 use crate::server::auth::WsRelayState;
2593 use crate::storage::HashtreeStore;
2594 use crate::test_support::{test_env_lock, EnvVarGuard};
2595 use axum::response::IntoResponse;
2596 use axum::{routing::post, Router};
2597 use base64::Engine;
2598 use hashtree_config::StorageBackend;
2599 use hashtree_core::sha256;
2600 use std::collections::HashSet;
2601 use std::sync::{Arc, Mutex as StdMutex};
2602 use std::time::Duration;
2603 use tempfile::TempDir;
2604
2605 fn test_app_state(store: Arc<HashtreeStore>) -> AppState {
2606 AppState {
2607 store,
2608 auth: None,
2609 daemon_started_at: 1_700_000_000,
2610 peer_mode: crate::config::ServerMode::Normal,
2611 hash_get_enabled: true,
2612 fips_endpoint: None,
2613 fips_blob_resolver: None,
2614 fetch_from_fips_peers: true,
2615 ws_relay: Arc::new(WsRelayState::new()),
2616 max_upload_bytes: 5 * 1024 * 1024,
2617 public_writes: true,
2618 public_plaintext_reads: true,
2619 require_random_untrusted_ingest: true,
2620 optimistic_blossom_uploads: false,
2621 optimistic_upload_queue_bytes: 256 * 1024 * 1024,
2622 optimistic_upload_queue: Arc::new(tokio::sync::Semaphore::new(256 * 1024 * 1024)),
2623 allowed_pubkeys: HashSet::new(),
2624 upstream_blossom: Vec::new(),
2625 upstream_http_client: super::super::new_upstream_http_client(),
2626 upstream_blossom_miss_cache: Arc::new(StdMutex::new(crate::server::new_lookup_cache())),
2627 upstream_blossom_fetch_metrics: Arc::new(
2628 crate::server::auth::UpstreamBlossomFetchMetrics::default(),
2629 ),
2630 blossom_upload_replicas: Vec::new(),
2631 blossom_upload_replica_queue_bytes: 256 * 1024 * 1024,
2632 blossom_upload_replica_queue: Arc::new(tokio::sync::Semaphore::new(256 * 1024 * 1024)),
2633 blossom_upload_replica_keys: None,
2634 blossom_upload_replica_scheduler: Arc::new(BlossomUploadReplicaScheduler::new()),
2635 social_graph: None,
2636 social_graph_store: None,
2637 social_graph_root: None,
2638 socialgraph_snapshot_public: false,
2639 nostr_relay: None,
2640 nostr_provider: None,
2641 nostr_relay_urls: Vec::new(),
2642 tree_root_cache: Arc::new(StdMutex::new(std::collections::HashMap::new())),
2643 inflight_blob_fetches: Arc::new(tokio::sync::Mutex::new(
2644 std::collections::HashMap::new(),
2645 )),
2646 inflight_blob_reads: Arc::new(
2647 tokio::sync::Mutex::new(std::collections::HashMap::new()),
2648 ),
2649 blob_cache: Arc::new(crate::blob_cache::BlobCache::for_tests()),
2650 directory_listing_cache: Arc::new(StdMutex::new(crate::server::new_lookup_cache())),
2651 resolved_path_cache: Arc::new(StdMutex::new(crate::server::new_lookup_cache())),
2652 thumbnail_path_cache: Arc::new(StdMutex::new(crate::server::new_lookup_cache())),
2653 cid_size_cache: Arc::new(StdMutex::new(crate::server::new_lookup_cache())),
2654 }
2655 }
2656
2657 async fn receive_replication<T>(receiver: &mut tokio::sync::mpsc::UnboundedReceiver<T>) -> T {
2658 tokio::time::timeout(Duration::from_secs(10), receiver.recv())
2659 .await
2660 .expect("replication request timed out")
2661 .expect("replication channel closed")
2662 }
2663
2664 fn create_upload_auth_header(keys: &nostr::Keys) -> String {
2665 use nostr::{EventBuilder, Kind, Tag, TagKind, Timestamp};
2666
2667 let now = Timestamp::now();
2668 let event = EventBuilder::new(Kind::Custom(BLOSSOM_AUTH_KIND), "")
2669 .tags(vec![
2670 Tag::custom(TagKind::Custom("t".into()), vec!["upload".to_string()]),
2671 Tag::custom(
2672 TagKind::Custom("expiration".into()),
2673 vec![(now.as_secs() + 300).to_string()],
2674 ),
2675 ])
2676 .custom_created_at(now)
2677 .sign_with_keys(keys)
2678 .expect("sign blossom auth");
2679 let json = serde_json::to_vec(&event).expect("serialize auth event");
2680 format!(
2681 "Nostr {}",
2682 base64::engine::general_purpose::STANDARD.encode(json)
2683 )
2684 }
2685
2686 fn create_batch_upload_auth_header(keys: &nostr::Keys, hashes: &[String]) -> String {
2687 use nostr::{EventBuilder, Kind, Tag, TagKind, Timestamp};
2688
2689 let now = Timestamp::now();
2690 let event = EventBuilder::new(Kind::Custom(BLOSSOM_AUTH_KIND), "")
2691 .tags(vec![
2692 Tag::custom(TagKind::Custom("t".into()), vec!["upload".to_string()]),
2693 Tag::custom(
2694 TagKind::Custom(BATCH_UPLOAD_HASH_LIST_AUTH_TAG.into()),
2695 vec![
2696 batch_upload_hash_list_digest(hashes.iter().map(String::as_str))
2697 .expect("batch hash list digest"),
2698 ],
2699 ),
2700 Tag::custom(
2701 TagKind::Custom("expiration".into()),
2702 vec![(now.as_secs() + 300).to_string()],
2703 ),
2704 ])
2705 .custom_created_at(now)
2706 .sign_with_keys(keys)
2707 .expect("sign blossom batch auth");
2708 let json = serde_json::to_vec(&event).expect("serialize batch auth event");
2709 format!(
2710 "Nostr {}",
2711 base64::engine::general_purpose::STANDARD.encode(json)
2712 )
2713 }
2714
2715 fn create_list_auth_header(keys: &nostr::Keys) -> String {
2716 use nostr::{EventBuilder, Kind, Tag, TagKind, Timestamp};
2717
2718 let now = Timestamp::now();
2719 let event = EventBuilder::new(Kind::Custom(BLOSSOM_AUTH_KIND), "")
2720 .tags(vec![
2721 Tag::custom(TagKind::Custom("t".into()), vec!["list".to_string()]),
2722 Tag::custom(
2723 TagKind::Custom("expiration".into()),
2724 vec![(now.as_secs() + 300).to_string()],
2725 ),
2726 ])
2727 .custom_created_at(now)
2728 .sign_with_keys(keys)
2729 .expect("sign blossom list auth");
2730 let json = serde_json::to_vec(&event).expect("serialize list auth event");
2731 format!(
2732 "Nostr {}",
2733 base64::engine::general_purpose::STANDARD.encode(json)
2734 )
2735 }
2736
2737 fn hosted_headers() -> HeaderMap {
2738 let mut headers = HeaderMap::new();
2739 headers.insert(header::HOST, "origin.internal".parse().unwrap());
2740 headers.insert("x-forwarded-host", "cdn.iris.to".parse().unwrap());
2741 headers.insert("x-forwarded-proto", "https".parse().unwrap());
2742 headers
2743 }
2744
2745 async fn read_descriptor(response: axum::response::Response) -> BlobDescriptor {
2746 let body = axum::body::to_bytes(response.into_body(), usize::MAX)
2747 .await
2748 .expect("read descriptor body");
2749 serde_json::from_slice(&body).expect("parse descriptor")
2750 }
2751
2752 fn upload_check_bits(response: UploadCheckResponse) -> Vec<bool> {
2753 let bytes = base64::engine::general_purpose::STANDARD
2754 .decode(response.present)
2755 .expect("decode upload check bitset");
2756 (0..response.count)
2757 .map(|index| bytes[index / 8] & (1 << (index % 8)) != 0)
2758 .collect()
2759 }
2760
2761 fn binary_batch_body(items: &[(&[u8], Option<&str>)]) -> Bytes {
2762 let mut body = Vec::new();
2763 body.extend_from_slice(BINARY_BATCH_UPLOAD_MAGIC);
2764 body.extend_from_slice(&(items.len() as u32).to_be_bytes());
2765 for (data, content_type) in items {
2766 body.extend_from_slice(&sha256(data));
2767 let content_type = content_type.unwrap_or("");
2768 body.extend_from_slice(&(content_type.len() as u16).to_be_bytes());
2769 body.extend_from_slice(&(*data).len().to_be_bytes());
2770 body.extend_from_slice(content_type.as_bytes());
2771 body.extend_from_slice(data);
2772 }
2773 Bytes::from(body)
2774 }
2775
2776 #[test]
2777 fn test_is_valid_sha256() {
2778 assert!(is_valid_sha256(
2779 "e2bab35b5296ec2242ded0a01f6d6723a5cd921239280c0a5f0b5589303336b6"
2780 ));
2781 assert!(is_valid_sha256(
2782 "0000000000000000000000000000000000000000000000000000000000000000"
2783 ));
2784
2785 assert!(!is_valid_sha256("e2bab35b5296ec2242ded0a01f6d6723"));
2787 assert!(!is_valid_sha256(
2789 "e2bab35b5296ec2242ded0a01f6d6723a5cd921239280c0a5f0b5589303336b6aa"
2790 ));
2791 assert!(!is_valid_sha256(
2793 "zzbab35b5296ec2242ded0a01f6d6723a5cd921239280c0a5f0b5589303336b6"
2794 ));
2795 assert!(!is_valid_sha256(""));
2797 }
2798
2799 #[test]
2800 fn test_parse_hash_and_extension() {
2801 let (hash, ext) = parse_hash_and_extension("abc123.png");
2802 assert_eq!(hash, "abc123");
2803 assert_eq!(ext, Some(".png"));
2804
2805 let (hash2, ext2) = parse_hash_and_extension("abc123");
2806 assert_eq!(hash2, "abc123");
2807 assert_eq!(ext2, None);
2808
2809 let (hash3, ext3) = parse_hash_and_extension("abc.123.jpg");
2810 assert_eq!(hash3, "abc.123");
2811 assert_eq!(ext3, Some(".jpg"));
2812 }
2813
2814 #[test]
2815 fn test_mime_to_extension() {
2816 assert_eq!(mime_to_extension("image/png"), ".png");
2817 assert_eq!(mime_to_extension("image/jpeg"), ".jpg");
2818 assert_eq!(mime_to_extension("video/mp4"), ".mp4");
2819 assert_eq!(mime_to_extension("application/octet-stream"), ".bin");
2820 assert_eq!(mime_to_extension("unknown/type"), ".bin");
2821 }
2822
2823 #[test]
2824 fn blossom_blob_url_uses_forwarded_public_origin_and_extension() {
2825 let headers = hosted_headers();
2826 let hash = "00".repeat(32);
2827
2828 assert_eq!(
2829 blossom_blob_url(&headers, &hash, "application/octet-stream"),
2830 format!("https://cdn.iris.to/{hash}.bin")
2831 );
2832 assert_eq!(
2833 blossom_blob_url(&headers, &hash, "image/png"),
2834 format!("https://cdn.iris.to/{hash}.png")
2835 );
2836 }
2837
2838 #[tokio::test]
2839 async fn upload_check_reports_present_hashes_in_request_order() {
2840 let temp_dir = TempDir::new().expect("temp dir");
2841 let store = Arc::new(
2842 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
2843 );
2844
2845 let present = b"present blob";
2846 let missing = b"missing blob";
2847 let present_hash = sha256(present);
2848 let missing_hash = sha256(missing);
2849 store.put_cached_blob(present).expect("seed blob");
2850
2851 let state = test_app_state(store);
2852 let response = upload_check(
2853 State(state),
2854 Json(UploadCheckRequest {
2855 hashes: vec![
2856 hex::encode(missing_hash),
2857 hex::encode(present_hash),
2858 hex::encode(present_hash),
2859 ],
2860 }),
2861 )
2862 .await
2863 .into_response();
2864
2865 assert_eq!(response.status(), StatusCode::OK);
2866 let body = axum::body::to_bytes(response.into_body(), usize::MAX)
2867 .await
2868 .expect("read response body");
2869 let parsed: UploadCheckResponse =
2870 serde_json::from_slice(&body).expect("parse upload check response");
2871 assert_eq!(upload_check_bits(parsed), vec![false, true, true]);
2872 }
2873
2874 #[tokio::test]
2875 async fn upload_check_rejects_invalid_hash() {
2876 let temp_dir = TempDir::new().expect("temp dir");
2877 let store = Arc::new(
2878 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
2879 );
2880 let state = test_app_state(store);
2881 let response = upload_check(
2882 State(state),
2883 Json(UploadCheckRequest {
2884 hashes: vec!["not-a-sha256".to_string()],
2885 }),
2886 )
2887 .await
2888 .into_response();
2889
2890 assert_eq!(response.status(), StatusCode::BAD_REQUEST);
2891 }
2892
2893 #[tokio::test]
2894 async fn upload_check_rejects_too_many_hashes() {
2895 let temp_dir = TempDir::new().expect("temp dir");
2896 let store = Arc::new(
2897 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
2898 );
2899 let state = test_app_state(store);
2900 let response = upload_check(
2901 State(state),
2902 Json(UploadCheckRequest {
2903 hashes: vec!["00".repeat(32); MAX_UPLOAD_CHECK_HASHES + 1],
2904 }),
2905 )
2906 .await
2907 .into_response();
2908
2909 assert_eq!(response.status(), StatusCode::PAYLOAD_TOO_LARGE);
2910 }
2911
2912 #[tokio::test]
2913 async fn upload_blob_batch_binary_replaces_cached_misses_after_commit() {
2914 let temp_dir = TempDir::new().expect("temp dir");
2915 let store = Arc::new(
2916 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
2917 );
2918 let state = test_app_state(Arc::clone(&store));
2919 let keys = nostr::Keys::generate();
2920 let mut headers = hosted_headers();
2921 headers.insert(
2922 header::AUTHORIZATION,
2923 create_upload_auth_header(&keys)
2924 .parse()
2925 .expect("auth header value"),
2926 );
2927 headers.insert(
2928 header::CONTENT_TYPE,
2929 "application/vnd.hashtree.blossom.batch.v1"
2930 .parse()
2931 .expect("content type header value"),
2932 );
2933
2934 let first = (0u8..=255).collect::<Vec<_>>();
2935 let second = (0u8..=255).map(|byte| byte ^ 0xaa).collect::<Vec<_>>();
2936 let first_hash = sha256(&first);
2937 let second_hash = sha256(&second);
2938 let body = binary_batch_body(&[
2939 (&first, Some("application/octet-stream")),
2940 (&second, Some("application/octet-stream")),
2941 ]);
2942
2943 let first_hash_hex = hex::encode(first_hash);
2944 let client = axum::extract::ConnectInfo("127.0.0.1:12345".parse().unwrap());
2945 let miss = head_blob(
2946 State(state.clone()),
2947 Path(format!("{first_hash_hex}.bin")),
2948 client,
2949 )
2950 .await
2951 .into_response();
2952 assert_eq!(miss.status(), StatusCode::NOT_FOUND);
2953
2954 let response = upload_blob_batch_binary(State(state.clone()), headers, body)
2955 .await
2956 .into_response();
2957
2958 assert_eq!(response.status(), StatusCode::OK);
2959 let body = axum::body::to_bytes(response.into_body(), usize::MAX)
2960 .await
2961 .expect("read batch response");
2962 let parsed: BatchUploadResponse =
2963 serde_json::from_slice(&body).expect("parse batch response");
2964 assert_eq!(parsed.uploaded, 2);
2965 assert_eq!(parsed.blobs.len(), 2);
2966 assert_eq!(parsed.blobs[0].sha256, hex::encode(first_hash));
2967 assert_eq!(parsed.blobs[1].sha256, hex::encode(second_hash));
2968 assert!(store.blob_exists(&first_hash).expect("first exists"));
2969 assert!(store.blob_exists(&second_hash).expect("second exists"));
2970
2971 let immediate_head = head_blob(State(state), Path(format!("{first_hash_hex}.bin")), client)
2972 .await
2973 .into_response();
2974 assert_eq!(
2975 immediate_head.status(),
2976 StatusCode::OK,
2977 "a committed batch write must replace a cached preflight miss",
2978 );
2979 }
2980
2981 #[tokio::test]
2982 async fn upload_blob_batch_binary_accepts_compact_batch_auth() {
2983 let temp_dir = TempDir::new().expect("temp dir");
2984 let store = Arc::new(
2985 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
2986 );
2987 let state = test_app_state(Arc::clone(&store));
2988 let keys = nostr::Keys::generate();
2989 let first = (0u8..=255).collect::<Vec<_>>();
2990 let second = (0u8..=255).map(|byte| byte ^ 0xaa).collect::<Vec<_>>();
2991 let first_hash = sha256(&first);
2992 let second_hash = sha256(&second);
2993 let hashes = vec![hex::encode(first_hash), hex::encode(second_hash)];
2994 let mut headers = hosted_headers();
2995 headers.insert(
2996 header::AUTHORIZATION,
2997 create_batch_upload_auth_header(&keys, &hashes)
2998 .parse()
2999 .expect("auth header value"),
3000 );
3001 headers.insert(
3002 header::CONTENT_TYPE,
3003 "application/vnd.hashtree.blossom.batch.v1"
3004 .parse()
3005 .expect("content type header value"),
3006 );
3007 let body = binary_batch_body(&[
3008 (&first, Some("application/octet-stream")),
3009 (&second, Some("application/octet-stream")),
3010 ]);
3011
3012 let response = upload_blob_batch_binary(State(state), headers, body)
3013 .await
3014 .into_response();
3015
3016 assert_eq!(response.status(), StatusCode::OK);
3017 assert!(store.blob_exists(&first_hash).expect("first exists"));
3018 assert!(store.blob_exists(&second_hash).expect("second exists"));
3019 }
3020
3021 #[tokio::test]
3022 async fn upload_blob_batch_binary_rejects_mismatched_compact_batch_auth() {
3023 let temp_dir = TempDir::new().expect("temp dir");
3024 let store = Arc::new(
3025 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
3026 );
3027 let state = test_app_state(Arc::clone(&store));
3028 let keys = nostr::Keys::generate();
3029 let first = (0u8..=255).collect::<Vec<_>>();
3030 let second = (0u8..=255).map(|byte| byte ^ 0xaa).collect::<Vec<_>>();
3031 let wrong_hashes = vec![hex::encode(sha256(&first)), "00".repeat(32)];
3032 let mut headers = hosted_headers();
3033 headers.insert(
3034 header::AUTHORIZATION,
3035 create_batch_upload_auth_header(&keys, &wrong_hashes)
3036 .parse()
3037 .expect("auth header value"),
3038 );
3039 headers.insert(
3040 header::CONTENT_TYPE,
3041 "application/vnd.hashtree.blossom.batch.v1"
3042 .parse()
3043 .expect("content type header value"),
3044 );
3045 let body = binary_batch_body(&[
3046 (&first, Some("application/octet-stream")),
3047 (&second, Some("application/octet-stream")),
3048 ]);
3049
3050 let response = upload_blob_batch_binary(State(state), headers, body)
3051 .await
3052 .into_response();
3053
3054 assert_eq!(response.status(), StatusCode::FORBIDDEN);
3055 assert!(!store.blob_exists(&sha256(&first)).expect("first absent"));
3056 assert!(!store.blob_exists(&sha256(&second)).expect("second absent"));
3057 }
3058
3059 #[tokio::test]
3060 async fn upload_blob_batch_rejects_missing_auth_before_decoding_payload() {
3061 let temp_dir = TempDir::new().expect("temp dir");
3062 let store = Arc::new(
3063 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
3064 );
3065 let state = test_app_state(store);
3066 let payload = BatchUploadRequest {
3067 blobs: vec![BatchUploadBlob {
3068 sha256: "00".repeat(32),
3069 content_type: None,
3070 data: "not-base64".to_string(),
3071 }],
3072 };
3073
3074 let response = upload_blob_batch(State(state), HeaderMap::new(), Json(payload))
3075 .await
3076 .into_response();
3077
3078 assert_eq!(response.status(), StatusCode::UNAUTHORIZED);
3079 }
3080
3081 #[tokio::test]
3082 async fn upload_blob_batch_binary_rejects_missing_auth_before_parsing_body() {
3083 let temp_dir = TempDir::new().expect("temp dir");
3084 let store = Arc::new(
3085 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
3086 );
3087 let state = test_app_state(store);
3088
3089 let response = upload_blob_batch_binary(
3090 State(state),
3091 HeaderMap::new(),
3092 Bytes::from_static(b"not a binary batch"),
3093 )
3094 .await
3095 .into_response();
3096
3097 assert_eq!(response.status(), StatusCode::UNAUTHORIZED);
3098 }
3099
3100 #[tokio::test(flavor = "current_thread")]
3101 async fn upload_blob_batch_binary_replicates_only_new_blobs() {
3102 let _lock = test_env_lock().lock().await;
3103 let config_dir = TempDir::new().expect("config dir");
3104 let _guard = EnvVarGuard::set("HTREE_CONFIG_DIR", config_dir.path());
3105
3106 let first = (0u8..=255).collect::<Vec<_>>();
3107 let second = (0u8..=255).map(|byte| byte ^ 0x55).collect::<Vec<_>>();
3108 let first_hash = sha256(&first);
3109 let second_hash = sha256(&second);
3110 let second_hash_hex = hex::encode(second_hash);
3111
3112 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<Vec<String>>();
3113 let replica_router = Router::new().route(
3114 "/upload/batch-binary",
3115 post(move |body: Bytes| {
3116 let tx = tx.clone();
3117 async move {
3118 let blobs = parse_binary_batch_upload(&body).expect("parse replica batch");
3119 let hashes = blobs
3120 .iter()
3121 .map(|blob| blob.sha256.clone())
3122 .collect::<Vec<_>>();
3123 let _ = tx.send(hashes.clone());
3124 Json(serde_json::json!({
3125 "uploaded": hashes.len(),
3126 "blobs": hashes.into_iter().map(|sha256| serde_json::json!({ "sha256": sha256 })).collect::<Vec<_>>(),
3127 }))
3128 }
3129 }),
3130 );
3131 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
3132 .await
3133 .expect("bind replica");
3134 let replica_addr = listener.local_addr().expect("replica addr");
3135 let _server_task =
3136 tokio::spawn(async move { axum::serve(listener, replica_router).await.unwrap() });
3137
3138 let temp_dir = TempDir::new().expect("temp dir");
3139 let store = Arc::new(
3140 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
3141 );
3142 store
3143 .put_cached_blobs(&[(first_hash, first.clone())])
3144 .expect("prestore first blob");
3145 let mut state = test_app_state(Arc::clone(&store));
3146 state.require_random_untrusted_ingest = false;
3147 state.blossom_upload_replicas = vec![format!("http://{replica_addr}")];
3148 state.blossom_upload_replica_keys = Some(Arc::new(nostr::Keys::generate()));
3149
3150 let keys = nostr::Keys::generate();
3151 let mut headers = hosted_headers();
3152 headers.insert(
3153 header::AUTHORIZATION,
3154 create_upload_auth_header(&keys)
3155 .parse()
3156 .expect("auth header value"),
3157 );
3158 headers.insert(
3159 header::CONTENT_TYPE,
3160 "application/vnd.hashtree.blossom.batch.v1"
3161 .parse()
3162 .expect("content type header value"),
3163 );
3164 let body = binary_batch_body(&[
3165 (&first, Some("application/octet-stream")),
3166 (&second, Some("application/octet-stream")),
3167 ]);
3168
3169 let response = upload_blob_batch_binary(State(state), headers, body)
3170 .await
3171 .into_response();
3172
3173 assert_eq!(response.status(), StatusCode::OK);
3174 let body = axum::body::to_bytes(response.into_body(), usize::MAX)
3175 .await
3176 .expect("read batch response");
3177 let parsed: BatchUploadResponse =
3178 serde_json::from_slice(&body).expect("parse batch response");
3179 assert_eq!(parsed.uploaded, 1);
3180 let replicated = receive_replication(&mut rx).await;
3181 assert_eq!(replicated, vec![second_hash_hex]);
3182 assert!(
3183 tokio::time::timeout(Duration::from_millis(100), rx.recv())
3184 .await
3185 .is_err(),
3186 "duplicate blob should not trigger a second replication batch"
3187 );
3188 }
3189
3190 #[tokio::test(flavor = "current_thread")]
3191 async fn upload_replication_coalesces_adjacent_binary_batches() {
3192 let _lock = test_env_lock().lock().await;
3193 let config_dir = TempDir::new().expect("config dir");
3194 let _guard = EnvVarGuard::set("HTREE_CONFIG_DIR", config_dir.path());
3195 let _flush_guard = EnvVarGuard::set("HTREE_BLOSSOM_REPLICA_COALESCE_FLUSH_MS", "2000");
3196 let _blobs_guard = EnvVarGuard::set("HTREE_BLOSSOM_REPLICA_COALESCE_MAX_BLOBS", "8");
3197 let _bytes_guard = EnvVarGuard::set("HTREE_BLOSSOM_REPLICA_COALESCE_MAX_BYTES", "1048576");
3198
3199 let first_a = b"coalesced-replication-first-a".to_vec();
3200 let first_b = b"coalesced-replication-first-b".to_vec();
3201 let second_a = b"coalesced-replication-second-a".to_vec();
3202 let second_b = b"coalesced-replication-second-b".to_vec();
3203 let expected_hashes = [&first_a, &first_b, &second_a, &second_b]
3204 .into_iter()
3205 .map(|data| hex::encode(sha256(data)))
3206 .collect::<HashSet<_>>();
3207
3208 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<Vec<String>>();
3209 let replica_router = Router::new().route(
3210 "/upload/batch-binary",
3211 post(move |body: Bytes| {
3212 let tx = tx.clone();
3213 async move {
3214 let blobs = parse_binary_batch_upload(&body).expect("parse replica batch");
3215 let hashes = blobs
3216 .iter()
3217 .map(|blob| blob.sha256.clone())
3218 .collect::<Vec<_>>();
3219 let _ = tx.send(hashes.clone());
3220 Json(serde_json::json!({
3221 "uploaded": hashes.len(),
3222 "blobs": hashes.into_iter().map(|sha256| serde_json::json!({ "sha256": sha256 })).collect::<Vec<_>>(),
3223 }))
3224 }
3225 }),
3226 );
3227 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
3228 .await
3229 .expect("bind replica");
3230 let replica_addr = listener.local_addr().expect("replica addr");
3231 let _server_task =
3232 tokio::spawn(async move { axum::serve(listener, replica_router).await.unwrap() });
3233
3234 let temp_dir = TempDir::new().expect("temp dir");
3235 let store = Arc::new(
3236 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
3237 );
3238 let mut state = test_app_state(store);
3239 state.require_random_untrusted_ingest = false;
3240 state.blossom_upload_replicas = vec![format!("http://{replica_addr}")];
3241 state.blossom_upload_replica_keys = Some(Arc::new(nostr::Keys::generate()));
3242 let metrics_before = blossom_upload_replica_queue_snapshot(&state);
3243
3244 let keys = nostr::Keys::generate();
3245 let mut headers = hosted_headers();
3246 headers.insert(
3247 header::AUTHORIZATION,
3248 create_upload_auth_header(&keys)
3249 .parse()
3250 .expect("auth header value"),
3251 );
3252 headers.insert(
3253 header::CONTENT_TYPE,
3254 "application/vnd.hashtree.blossom.batch.v1"
3255 .parse()
3256 .expect("content type header value"),
3257 );
3258 let first_body = binary_batch_body(&[
3259 (&first_a, Some("application/octet-stream")),
3260 (&first_b, Some("application/octet-stream")),
3261 ]);
3262 let second_body = binary_batch_body(&[
3263 (&second_a, Some("application/octet-stream")),
3264 (&second_b, Some("application/octet-stream")),
3265 ]);
3266
3267 let first_response =
3268 upload_blob_batch_binary(State(state.clone()), headers.clone(), first_body)
3269 .await
3270 .into_response();
3271 assert_eq!(first_response.status(), StatusCode::OK);
3272 let second_response = upload_blob_batch_binary(State(state.clone()), headers, second_body)
3273 .await
3274 .into_response();
3275 assert_eq!(second_response.status(), StatusCode::OK);
3276
3277 let replicated = receive_replication(&mut rx).await;
3278 let replicated_hashes = replicated.into_iter().collect::<HashSet<_>>();
3279 assert_eq!(replicated_hashes, expected_hashes);
3280 assert!(
3281 tokio::time::timeout(Duration::from_millis(150), rx.recv())
3282 .await
3283 .is_err(),
3284 "adjacent batches should be merged into one replica request"
3285 );
3286 let metrics_after = blossom_upload_replica_queue_snapshot(&state);
3287 assert!(
3288 metrics_after.accepted_batches > metrics_before.accepted_batches,
3289 "coalesced replication should increment accepted batch metrics"
3290 );
3291 assert!(
3292 metrics_after.accepted_blobs >= metrics_before.accepted_blobs + 4,
3293 "coalesced replication should increment accepted blob metrics"
3294 );
3295 }
3296
3297 #[test]
3298 fn binary_batch_parser_rejects_trailing_bytes() {
3299 let data = (0u8..=255).collect::<Vec<_>>();
3300 let mut body = binary_batch_body(&[(&data, None)]).to_vec();
3301 body.push(0);
3302
3303 let error = parse_binary_batch_upload(&body).expect_err("trailing bytes rejected");
3304
3305 assert_eq!(error.0, StatusCode::BAD_REQUEST);
3306 assert!(error.1.contains("trailing"));
3307 }
3308
3309 #[tokio::test]
3310 async fn head_upload_accepts_valid_bud06_preflight() {
3311 let temp_dir = TempDir::new().expect("temp dir");
3312 let store = Arc::new(
3313 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
3314 );
3315 let state = test_app_state(store);
3316 let mut headers = hosted_headers();
3317 headers.insert("x-sha-256", "00".repeat(32).parse().unwrap());
3318 headers.insert("x-content-length", "16".parse().unwrap());
3319 headers.insert(
3320 "x-content-type",
3321 "application/octet-stream".parse().unwrap(),
3322 );
3323
3324 let response = head_upload(State(state), headers).await.into_response();
3325
3326 assert_eq!(response.status(), StatusCode::OK);
3327 }
3328
3329 #[tokio::test]
3330 async fn upload_blob_returns_bud02_statuses_and_public_descriptor_url() {
3331 let temp_dir = TempDir::new().expect("temp dir");
3332 let store = Arc::new(
3333 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
3334 );
3335 let state = test_app_state(store);
3336 let keys = nostr::Keys::generate();
3337 let mut headers = hosted_headers();
3338 headers.insert(
3339 header::AUTHORIZATION,
3340 create_upload_auth_header(&keys)
3341 .parse()
3342 .expect("auth header value"),
3343 );
3344 headers.insert(
3345 header::CONTENT_TYPE,
3346 "application/octet-stream"
3347 .parse()
3348 .expect("content type header value"),
3349 );
3350
3351 let body = axum::body::Bytes::from((0u8..=255).collect::<Vec<_>>());
3352 let hash_hex = hex::encode(sha256(&body));
3353 let first = upload_blob(State(state.clone()), headers.clone(), body.clone())
3354 .await
3355 .into_response();
3356 assert_eq!(first.status(), StatusCode::CREATED);
3357 let first_descriptor = read_descriptor(first).await;
3358 assert_eq!(
3359 first_descriptor.url,
3360 format!("https://cdn.iris.to/{hash_hex}.bin")
3361 );
3362 assert_eq!(first_descriptor.sha256, hash_hex);
3363
3364 let second = upload_blob(State(state), headers, body)
3365 .await
3366 .into_response();
3367 assert_eq!(second.status(), StatusCode::OK);
3368 let second_descriptor = read_descriptor(second).await;
3369 assert_eq!(second_descriptor.url, first_descriptor.url);
3370 }
3371
3372 #[tokio::test(flavor = "current_thread")]
3373 async fn upload_blob_replicates_to_configured_blossom_target() {
3374 let _lock = test_env_lock().lock().await;
3375 let config_dir = TempDir::new().expect("config dir");
3376 let _guard = EnvVarGuard::set("HTREE_CONFIG_DIR", config_dir.path());
3377
3378 let data = Bytes::from_static(b"write-behind-replication-data");
3379 let expected_hash = hex::encode(sha256(data.as_ref()));
3380 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<usize>();
3381 let response_hash = expected_hash.clone();
3382 let replica_router = Router::new().route(
3383 "/upload/batch-binary",
3384 post(move |body: Bytes| {
3385 let tx = tx.clone();
3386 let response_hash = response_hash.clone();
3387 async move {
3388 let _ = tx.send(body.len());
3389 Json(serde_json::json!({
3390 "uploaded": 1,
3391 "blobs": [{"sha256": response_hash}],
3392 }))
3393 }
3394 }),
3395 );
3396 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
3397 .await
3398 .expect("bind replica");
3399 let replica_addr = listener.local_addr().expect("replica addr");
3400 let _server_task =
3401 tokio::spawn(async move { axum::serve(listener, replica_router).await.unwrap() });
3402
3403 let temp = TempDir::new().expect("tempdir");
3404 let store = Arc::new(HashtreeStore::new(temp.path()).expect("store"));
3405 let mut state = test_app_state(store);
3406 state.require_random_untrusted_ingest = false;
3407 state.blossom_upload_replicas = vec![format!("http://{replica_addr}")];
3408 state.blossom_upload_replica_keys = Some(Arc::new(nostr::Keys::generate()));
3409
3410 let keys = nostr::Keys::generate();
3411 let mut headers = HeaderMap::new();
3412 headers.insert(
3413 header::AUTHORIZATION,
3414 create_upload_auth_header(&keys).parse().unwrap(),
3415 );
3416 headers.insert(
3417 header::CONTENT_TYPE,
3418 "application/octet-stream".parse().unwrap(),
3419 );
3420
3421 let response = upload_blob(State(state), headers, data)
3422 .await
3423 .into_response();
3424 assert_eq!(response.status(), StatusCode::CREATED);
3425 let replicated = receive_replication(&mut rx).await;
3426 assert!(replicated > 0);
3427 assert_eq!(expected_hash.len(), 64);
3428 }
3429
3430 #[tokio::test(flavor = "current_thread")]
3431 async fn upload_blob_duplicate_does_not_replicate_to_configured_blossom_target() {
3432 let _lock = test_env_lock().lock().await;
3433 let config_dir = TempDir::new().expect("config dir");
3434 let _guard = EnvVarGuard::set("HTREE_CONFIG_DIR", config_dir.path());
3435
3436 let data = Bytes::from_static(b"write-behind-duplicate-raw-data");
3437 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<usize>();
3438 let response_hash = hex::encode(sha256(data.as_ref()));
3439 let replica_router = Router::new().route(
3440 "/upload/batch-binary",
3441 post(move |body: Bytes| {
3442 let tx = tx.clone();
3443 let response_hash = response_hash.clone();
3444 async move {
3445 let _ = tx.send(body.len());
3446 Json(serde_json::json!({
3447 "uploaded": 1,
3448 "blobs": [{"sha256": response_hash}],
3449 }))
3450 }
3451 }),
3452 );
3453 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
3454 .await
3455 .expect("bind replica");
3456 let replica_addr = listener.local_addr().expect("replica addr");
3457 let _server_task =
3458 tokio::spawn(async move { axum::serve(listener, replica_router).await.unwrap() });
3459
3460 let temp = TempDir::new().expect("tempdir");
3461 let store = Arc::new(HashtreeStore::new(temp.path()).expect("store"));
3462 store.put_cached_blob(&data).expect("seed duplicate blob");
3463 let mut state = test_app_state(store);
3464 state.require_random_untrusted_ingest = false;
3465 state.blossom_upload_replicas = vec![format!("http://{replica_addr}")];
3466 state.blossom_upload_replica_keys = Some(Arc::new(nostr::Keys::generate()));
3467
3468 let keys = nostr::Keys::generate();
3469 let mut headers = HeaderMap::new();
3470 headers.insert(
3471 header::AUTHORIZATION,
3472 create_upload_auth_header(&keys).parse().unwrap(),
3473 );
3474 headers.insert(
3475 header::CONTENT_TYPE,
3476 "application/octet-stream".parse().unwrap(),
3477 );
3478
3479 let response = upload_blob(State(state), headers, data)
3480 .await
3481 .into_response();
3482 assert_eq!(response.status(), StatusCode::OK);
3483 assert!(
3484 tokio::time::timeout(Duration::from_millis(100), rx.recv())
3485 .await
3486 .is_err(),
3487 "duplicate raw upload should not trigger write-behind replication"
3488 );
3489 }
3490
3491 #[tokio::test]
3492 async fn list_blobs_returns_public_descriptor_urls_with_extensions() {
3493 let temp_dir = TempDir::new().expect("temp dir");
3494 let store = Arc::new(
3495 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
3496 );
3497 let keys = nostr::Keys::generate();
3498 let pubkey_hex = keys.public_key().to_hex();
3499 let pubkey_bytes: [u8; 32] = from_hex(&pubkey_hex).expect("pubkey bytes");
3500 let body = (0u8..=255).collect::<Vec<_>>();
3501 let hash_hex = store
3502 .put_owned_blob(&body, &pubkey_bytes)
3503 .expect("store owned blob");
3504 let state = test_app_state(store);
3505 let mut headers = hosted_headers();
3506 headers.insert(
3507 header::AUTHORIZATION,
3508 create_list_auth_header(&keys)
3509 .parse()
3510 .expect("auth header value"),
3511 );
3512
3513 let response = list_blobs(
3514 State(state),
3515 Path(pubkey_hex),
3516 Query(ListQuery {
3517 since: None,
3518 until: None,
3519 limit: None,
3520 cursor: None,
3521 }),
3522 headers,
3523 )
3524 .await
3525 .into_response();
3526
3527 assert_eq!(response.status(), StatusCode::OK);
3528 let body = axum::body::to_bytes(response.into_body(), usize::MAX)
3529 .await
3530 .expect("read list body");
3531 let descriptors: Vec<BlobDescriptor> =
3532 serde_json::from_slice(&body).expect("parse descriptor list");
3533 assert_eq!(descriptors.len(), 1);
3534 assert_eq!(
3535 descriptors[0].url,
3536 format!("https://cdn.iris.to/{hash_hex}.bin")
3537 );
3538 }
3539
3540 #[tokio::test]
3541 async fn optimistic_uploads_return_accepted_and_store_in_background() {
3542 let temp_dir = TempDir::new().expect("temp dir");
3543 let store = Arc::new(
3544 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
3545 );
3546 let mut state = test_app_state(Arc::clone(&store));
3547 state.optimistic_blossom_uploads = true;
3548
3549 let keys = nostr::Keys::generate();
3550 let mut headers = HeaderMap::new();
3551 headers.insert(
3552 header::AUTHORIZATION,
3553 create_upload_auth_header(&keys)
3554 .parse()
3555 .expect("auth header value"),
3556 );
3557 headers.insert(
3558 header::CONTENT_TYPE,
3559 "application/octet-stream"
3560 .parse()
3561 .expect("content type header value"),
3562 );
3563
3564 let body = axum::body::Bytes::from((0u8..=255).map(|byte| byte ^ 0x55).collect::<Vec<_>>());
3565 let hash = sha256(&body);
3566 let response = upload_blob(State(state), headers, body)
3567 .await
3568 .into_response();
3569 assert_eq!(response.status(), StatusCode::ACCEPTED);
3570
3571 for _ in 0..50 {
3572 if store.blob_exists(&hash).expect("blob exists check") {
3573 return;
3574 }
3575 tokio::time::sleep(Duration::from_millis(10)).await;
3576 }
3577
3578 panic!("optimistic upload was not stored in the background");
3579 }
3580
3581 #[tokio::test]
3582 async fn optimistic_upload_existing_blob_skips_queue() {
3583 let temp_dir = TempDir::new().expect("temp dir");
3584 let store = Arc::new(
3585 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
3586 );
3587 let body = axum::body::Bytes::from(
3588 (0u16..=255)
3589 .map(|value| ((value * 73 + 19) % 256) as u8)
3590 .collect::<Vec<_>>(),
3591 );
3592 store.put_cached_blob(&body).expect("seed blob");
3593
3594 let mut state = test_app_state(Arc::clone(&store));
3595 state.optimistic_blossom_uploads = true;
3596 state.optimistic_upload_queue_bytes = 1;
3597 state.optimistic_upload_queue = Arc::new(tokio::sync::Semaphore::new(1));
3598
3599 let keys = nostr::Keys::generate();
3600 let mut headers = HeaderMap::new();
3601 headers.insert(
3602 header::AUTHORIZATION,
3603 create_upload_auth_header(&keys)
3604 .parse()
3605 .expect("auth header value"),
3606 );
3607 headers.insert(
3608 header::CONTENT_TYPE,
3609 "application/octet-stream"
3610 .parse()
3611 .expect("content type header value"),
3612 );
3613
3614 let response = upload_blob(State(state), headers, body)
3615 .await
3616 .into_response();
3617 assert_eq!(response.status(), StatusCode::OK);
3618 }
3619
3620 #[tokio::test]
3621 async fn optimistic_upload_existing_blob_uses_queue_before_preflight_when_queue_has_room() {
3622 let temp_dir = TempDir::new().expect("temp dir");
3623 let store = Arc::new(
3624 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
3625 );
3626 let body = axum::body::Bytes::from((0u8..=255).rev().collect::<Vec<_>>());
3627 let hash_hex = hex::encode(sha256(&body));
3628 store.put_cached_blob(&body).expect("seed blob");
3629
3630 let mut state = test_app_state(store);
3631 state.optimistic_blossom_uploads = true;
3632
3633 let keys = nostr::Keys::generate();
3634 let mut headers = HeaderMap::new();
3635 headers.insert(
3636 header::AUTHORIZATION,
3637 create_upload_auth_header(&keys)
3638 .parse()
3639 .expect("auth header value"),
3640 );
3641 headers.insert(
3642 header::CONTENT_TYPE,
3643 "application/octet-stream"
3644 .parse()
3645 .expect("content type header value"),
3646 );
3647
3648 let response = upload_blob(State(state), headers, body)
3649 .await
3650 .into_response();
3651 assert_eq!(response.status(), StatusCode::ACCEPTED);
3652
3653 for _ in 0..50 {
3654 if !optimistic_upload_is_inflight(&hash_hex) {
3655 return;
3656 }
3657 tokio::time::sleep(Duration::from_millis(10)).await;
3658 }
3659
3660 clear_optimistic_upload_inflight(&hash_hex);
3661 panic!("optimistic upload in-flight marker was not cleared");
3662 }
3663
3664 #[tokio::test]
3665 async fn optimistic_upload_existing_blob_does_not_replicate_duplicate() {
3666 let _lock = test_env_lock().lock().await;
3667 let config_dir = TempDir::new().expect("config dir");
3668 let _guard = EnvVarGuard::set("HTREE_CONFIG_DIR", config_dir.path());
3669
3670 let temp_dir = TempDir::new().expect("temp dir");
3671 let store = Arc::new(
3672 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
3673 );
3674 let body = axum::body::Bytes::from((0u8..=255).map(|byte| byte ^ 0xaa).collect::<Vec<_>>());
3675 let hash_hex = hex::encode(sha256(&body));
3676 store.put_cached_blob(&body).expect("seed blob");
3677
3678 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<usize>();
3679 let response_hash = hash_hex.clone();
3680 let replica_router = Router::new().route(
3681 "/upload/batch-binary",
3682 post(move |body: Bytes| {
3683 let tx = tx.clone();
3684 let response_hash = response_hash.clone();
3685 async move {
3686 let _ = tx.send(body.len());
3687 Json(serde_json::json!({
3688 "uploaded": 1,
3689 "blobs": [{"sha256": response_hash}],
3690 }))
3691 }
3692 }),
3693 );
3694 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
3695 .await
3696 .expect("bind replica");
3697 let replica_addr = listener.local_addr().expect("replica addr");
3698 let _server_task =
3699 tokio::spawn(async move { axum::serve(listener, replica_router).await.unwrap() });
3700
3701 let mut state = test_app_state(store);
3702 state.optimistic_blossom_uploads = true;
3703 state.require_random_untrusted_ingest = false;
3704 state.blossom_upload_replicas = vec![format!("http://{replica_addr}")];
3705 state.blossom_upload_replica_keys = Some(Arc::new(nostr::Keys::generate()));
3706
3707 let keys = nostr::Keys::generate();
3708 let mut headers = HeaderMap::new();
3709 headers.insert(
3710 header::AUTHORIZATION,
3711 create_upload_auth_header(&keys)
3712 .parse()
3713 .expect("auth header value"),
3714 );
3715 headers.insert(
3716 header::CONTENT_TYPE,
3717 "application/octet-stream"
3718 .parse()
3719 .expect("content type header value"),
3720 );
3721
3722 let response = upload_blob(State(state), headers, body)
3723 .await
3724 .into_response();
3725 assert_eq!(response.status(), StatusCode::ACCEPTED);
3726
3727 for _ in 0..50 {
3728 if !optimistic_upload_is_inflight(&hash_hex) {
3729 break;
3730 }
3731 tokio::time::sleep(Duration::from_millis(10)).await;
3732 }
3733 clear_optimistic_upload_inflight(&hash_hex);
3734
3735 assert!(
3736 tokio::time::timeout(Duration::from_millis(100), rx.recv())
3737 .await
3738 .is_err(),
3739 "optimistic duplicate upload should not trigger write-behind replication"
3740 );
3741 }
3742
3743 #[tokio::test]
3744 async fn optimistic_upload_inflight_duplicate_skips_queue() {
3745 let temp_dir = TempDir::new().expect("temp dir");
3746 let store = Arc::new(
3747 HashtreeStore::with_options(temp_dir.path(), None, 128 * 1024 * 1024).expect("store"),
3748 );
3749 let body = axum::body::Bytes::from((0u8..=255).collect::<Vec<_>>());
3750 let hash_hex = hex::encode(sha256(&body));
3751 assert!(mark_optimistic_upload_inflight(&hash_hex));
3752
3753 let mut state = test_app_state(store);
3754 state.optimistic_blossom_uploads = true;
3755 state.optimistic_upload_queue_bytes = 1;
3756 state.optimistic_upload_queue = Arc::new(tokio::sync::Semaphore::new(1));
3757
3758 let keys = nostr::Keys::generate();
3759 let mut headers = HeaderMap::new();
3760 headers.insert(
3761 header::AUTHORIZATION,
3762 create_upload_auth_header(&keys)
3763 .parse()
3764 .expect("auth header value"),
3765 );
3766 headers.insert(
3767 header::CONTENT_TYPE,
3768 "application/octet-stream"
3769 .parse()
3770 .expect("content type header value"),
3771 );
3772
3773 let response = upload_blob(State(state), headers, body)
3774 .await
3775 .into_response();
3776 clear_optimistic_upload_inflight(&hash_hex);
3777 assert_eq!(response.status(), StatusCode::ACCEPTED);
3778 }
3779
3780 #[test]
3781 fn public_writes_accept_unlisted_authors_for_uploads() {
3782 let temp_dir = TempDir::new().expect("temp dir");
3783 let store =
3784 Arc::new(HashtreeStore::with_options(temp_dir.path(), None, 700).expect("store"));
3785 let mut state = test_app_state(store);
3786 let pubkey = "ea4fe79e57f209309bffed2f92f0b95b59d3d1cb4e8892444398aeea7ee317ed";
3787
3788 state.public_writes = true;
3789 assert!(can_accept_upload_author(&state, pubkey));
3790 assert!(!is_allowed_write_author(&state, pubkey));
3791
3792 state.public_writes = false;
3793 assert!(!can_accept_upload_author(&state, pubkey));
3794 }
3795
3796 #[test]
3797 fn public_write_trust_allows_octet_stream_and_raw_media_payloads() {
3798 let encrypted_block: Vec<u8> = (0..=255).collect();
3799
3800 assert_eq!(
3801 validate_upload_payload(&encrypted_block, "application/octet-stream", false, true,),
3802 Ok(())
3803 );
3804
3805 assert_eq!(
3806 validate_upload_payload(b"audio bytes", "audio/mpeg", true, true,),
3807 Ok(())
3808 );
3809
3810 assert_eq!(
3811 validate_upload_payload(b"audio bytes", "audio/mpeg", false, true,),
3812 Err((
3813 StatusCode::FORBIDDEN,
3814 "Raw media uploads require write access".to_string(),
3815 ))
3816 );
3817 }
3818
3819 #[test]
3820 fn authenticated_chk_uploads_skip_entropy_heuristic() {
3821 let low_unique_block: Vec<u8> = (0..256).map(|i| (i % 139) as u8).collect();
3822
3823 assert_eq!(
3824 validate_upload_payload(&low_unique_block, "application/octet-stream", true, true,),
3825 Ok(())
3826 );
3827
3828 assert_eq!(
3829 validate_upload_payload(&low_unique_block, "application/octet-stream", false, true,),
3830 Err((
3831 StatusCode::UNSUPPORTED_MEDIA_TYPE,
3832 "Data not encrypted. Unique: 139 (min: 140)".to_string(),
3833 ))
3834 );
3835 }
3836
3837 #[test]
3838 fn unowned_public_uploads_use_cache_storage_semantics() {
3839 let temp_dir = TempDir::new().expect("temp dir");
3840 let store = Arc::new(
3841 HashtreeStore::with_options_and_backend(
3842 temp_dir.path(),
3843 None,
3844 700,
3845 true,
3846 &StorageBackend::Fs,
3847 )
3848 .expect("store"),
3849 );
3850 let state = test_app_state(Arc::clone(&store));
3851
3852 let owned = vec![1u8; 280];
3853 let owned_hash = sha256(&owned);
3854 store_blossom_blob(&state, &owned, &owned_hash, &[2u8; 32], true).expect("owned upload");
3855
3856 let public_upload = vec![3u8; 280];
3857 let public_hash = sha256(&public_upload);
3858 store_blossom_blob(&state, &public_upload, &public_hash, &[4u8; 32], false)
3859 .expect("public upload");
3860
3861 let replacement = vec![5u8; 280];
3862 let replacement_hash = sha256(&replacement);
3863 state
3864 .store
3865 .put_cached_blob(&replacement)
3866 .expect("replacement cached blob");
3867
3868 assert!(state.store.blob_exists(&owned_hash).expect("owned exists"));
3869 assert!(!state
3870 .store
3871 .blob_exists(&public_hash)
3872 .expect("public upload evicted"));
3873 assert!(state
3874 .store
3875 .blob_exists(&replacement_hash)
3876 .expect("replacement exists"));
3877 assert!(state
3878 .store
3879 .is_blob_owner(&owned_hash, &[2u8; 32])
3880 .expect("owned tracked"));
3881 assert!(!state
3882 .store
3883 .blob_has_owners(&public_hash)
3884 .expect("public upload unowned"));
3885 }
3886
3887 #[test]
3888 fn owned_blossom_uploads_are_rejected_when_storage_limit_is_full() {
3889 let temp_dir = TempDir::new().expect("temp dir");
3890 let store = Arc::new(
3891 HashtreeStore::with_options_and_backend(
3892 temp_dir.path(),
3893 None,
3894 500,
3895 true,
3896 &StorageBackend::Fs,
3897 )
3898 .expect("store"),
3899 );
3900 let state = test_app_state(Arc::clone(&store));
3901
3902 let first = vec![1u8; 300];
3903 let first_hash = sha256(&first);
3904 let owner = [2u8; 32];
3905 store_blossom_blob(&state, &first, &first_hash, &owner, true).expect("first upload");
3906
3907 let second = vec![3u8; 300];
3908 let second_hash = sha256(&second);
3909 let error = store_blossom_blob(&state, &second, &second_hash, &owner, true)
3910 .expect_err("second owned upload should exceed the storage limit");
3911
3912 assert!(
3913 error.to_string().contains("storage limit"),
3914 "unexpected error: {error}"
3915 );
3916 assert!(state
3917 .store
3918 .blob_exists(&first_hash)
3919 .expect("first blob remains"));
3920 assert!(!state
3921 .store
3922 .blob_exists(&second_hash)
3923 .expect("second blob rejected"));
3924 assert!(state
3925 .store
3926 .is_blob_owner(&first_hash, &owner)
3927 .expect("first owner tracked"));
3928 assert!(!state
3929 .store
3930 .is_blob_owner(&second_hash, &owner)
3931 .expect("second owner not tracked"));
3932 }
3933}