1use std::path::Path;
2
3use axum::body::Bytes;
4use shardline_protocol::RepositoryScope;
5use shardline_storage::{ObjectBody, ObjectIntegrity, ObjectKey, ObjectStore, PutOutcome};
6
7use crate::{
8 ServerError, ShardMetadataLimits,
9 model::UploadFileResponse,
10 protocol_support::shared_sha256_object_key,
11 upload_ingest::{FileUploadIngestor, RequestBodyReader, read_body_to_bytes},
12 validation::validate_identifier,
13 xet_adapter::{
14 ShardUploadResponse, XorbUploadResponse, register_uploaded_shard_bytes,
15 store_uploaded_xorb_bytes,
16 },
17};
18
19impl super::PostgresBackend {
20 pub async fn upload_file(
27 &self,
28 file_id: &str,
29 body: Bytes,
30 repository_scope: Option<&RepositoryScope>,
31 ) -> Result<UploadFileResponse, ServerError> {
32 self.upload_file_stream(
33 file_id,
34 RequestBodyReader::from_bytes(body),
35 repository_scope,
36 None,
37 )
38 .await
39 }
40
41 pub(crate) async fn upload_file_stream(
48 &self,
49 file_id: &str,
50 mut body: RequestBodyReader,
51 repository_scope: Option<&RepositoryScope>,
52 expected_sha256: Option<&str>,
53 ) -> Result<UploadFileResponse, ServerError> {
54 validate_identifier(file_id)?;
55
56 let object_store = self.object_store();
57 let mut ingestor = FileUploadIngestor::new_with_parallelism(
58 self.chunk_size,
59 expected_sha256.is_some(),
60 self.upload_max_in_flight_chunks,
61 );
62 while let Some(bytes) = body.next_bytes().await? {
63 ingestor.ingest_body_chunk(&object_store, &bytes).await?;
64 }
65
66 let (record, response) = ingestor
67 .finish(&object_store, file_id, repository_scope, expected_sha256)
68 .await?;
69 self.record_store
70 .commit_file_version_metadata(&record)
71 .await?;
72
73 Ok(response)
74 }
75
76 pub(crate) fn put_object_bytes_if_absent(
77 &self,
78 object_key: &ObjectKey,
79 bytes: Vec<u8>,
80 ) -> Result<PutOutcome, ServerError> {
81 let integrity = ObjectIntegrity::new(
82 shardline_protocol::ShardlineHash::from_bytes(*blake3::hash(&bytes).as_bytes()),
83 u64::try_from(bytes.len())?,
84 );
85 Ok(self.object_store().put_if_absent(
86 object_key,
87 ObjectBody::from_vec(bytes),
88 &integrity,
89 )?)
90 }
91
92 pub(crate) fn put_sha256_addressed_object_bytes_if_absent(
93 &self,
94 object_key: &ObjectKey,
95 digest_hex: &str,
96 bytes: Vec<u8>,
97 ) -> Result<PutOutcome, ServerError> {
98 let canonical_key = shared_sha256_object_key(digest_hex)?;
99 let integrity = ObjectIntegrity::new(
100 shardline_protocol::ShardlineHash::from_bytes(*blake3::hash(&bytes).as_bytes()),
101 u64::try_from(bytes.len())?,
102 );
103 let canonical_outcome = self.object_store().put_if_absent(
104 &canonical_key,
105 ObjectBody::from_vec(bytes),
106 &integrity,
107 )?;
108 if canonical_key == *object_key {
109 return Ok(canonical_outcome);
110 }
111 Ok(self
112 .object_store()
113 .copy_if_absent(&canonical_key, object_key)?)
114 }
115
116 pub(crate) fn copy_object_if_absent(
117 &self,
118 source: &ObjectKey,
119 destination: &ObjectKey,
120 ) -> Result<PutOutcome, ServerError> {
121 Ok(self.object_store().copy_if_absent(source, destination)?)
122 }
123
124 pub(crate) fn put_object_bytes_overwrite(
125 &self,
126 object_key: &ObjectKey,
127 bytes: Vec<u8>,
128 ) -> Result<(), ServerError> {
129 let integrity = ObjectIntegrity::new(
130 shardline_protocol::ShardlineHash::from_bytes(*blake3::hash(&bytes).as_bytes()),
131 u64::try_from(bytes.len())?,
132 );
133 Ok(self.object_store().put_overwrite(
134 object_key,
135 ObjectBody::from_vec(bytes),
136 &integrity,
137 )?)
138 }
139
140 pub(crate) fn put_sha256_addressed_object_file(
141 &self,
142 object_key: &ObjectKey,
143 digest_hex: &str,
144 path: &Path,
145 integrity: &ObjectIntegrity,
146 ) -> Result<PutOutcome, ServerError> {
147 let canonical_key = shared_sha256_object_key(digest_hex)?;
148 let canonical_outcome =
149 self.object_store()
150 .put_content_addressed_file(&canonical_key, path, integrity)?;
151 if canonical_key == *object_key {
152 return Ok(canonical_outcome);
153 }
154 Ok(self
155 .object_store()
156 .copy_if_absent(&canonical_key, object_key)?)
157 }
158
159 pub async fn upload_xorb(
166 &self,
167 expected_hash: &str,
168 body: Bytes,
169 ) -> Result<XorbUploadResponse, ServerError> {
170 self.upload_xorb_stream(expected_hash, RequestBodyReader::from_bytes(body))
171 .await
172 }
173
174 pub(crate) async fn upload_xorb_stream(
181 &self,
182 expected_hash: &str,
183 mut body: RequestBodyReader,
184 ) -> Result<XorbUploadResponse, ServerError> {
185 let uploaded_body = read_body_to_bytes(&mut body).await?;
186 let object_store = self.object_store();
187 store_uploaded_xorb_bytes(&object_store, expected_hash, &uploaded_body)
188 .map_err(ServerError::from)
189 }
190
191 pub(crate) async fn upload_shard_stream(
198 &self,
199 mut body: RequestBodyReader,
200 repository_scope: Option<&RepositoryScope>,
201 shard_metadata_limits: ShardMetadataLimits,
202 ) -> Result<ShardUploadResponse, ServerError> {
203 let uploaded_body = read_body_to_bytes(&mut body).await?;
204 let record_store = self.record_store.clone();
205 let object_store = self.object_store();
206 register_uploaded_shard_bytes(
207 &object_store,
208 &uploaded_body,
209 repository_scope,
210 shard_metadata_limits,
211 move |records, mappings| async move {
212 record_store
213 .commit_native_shard_metadata(&records, &mappings)
214 .await?;
215 Ok(())
216 },
217 )
218 .await
219 .map_err(ServerError::from)
220 }
221}
222
223#[cfg(test)]
224mod tests {
225 use std::num::NonZeroUsize;
226
227 use sha2::Digest;
228 use shardline_storage::ObjectKey;
229
230 use super::super::PostgresBackend;
231 use super::*;
232 use crate::object_store::ServerObjectStore;
233 use crate::protocol_support::shared_sha256_object_key;
234 use crate::test_fixtures::single_chunk_xorb;
235 use crate::upload_ingest::RequestBodyReader;
236
237 const TEST_PG_URL: &str = "postgres://localhost:5432/test";
238
239 fn make_object_key(label: &str) -> ObjectKey {
240 ObjectKey::parse(&format!("test/{label}")).unwrap()
241 }
242
243 async fn make_backend() -> (PostgresBackend, tempfile::TempDir) {
244 let root = tempfile::tempdir().expect("temp dir");
245 let object_store =
246 ServerObjectStore::local(root.path().join("chunks")).expect("local store");
247 let backend = PostgresBackend::new_with_object_store_and_upload_parallelism(
248 root.path().to_path_buf(),
249 "http://127.0.0.1:8080".to_owned(),
250 NonZeroUsize::new(65536).unwrap(),
251 NonZeroUsize::new(64).unwrap(),
252 TEST_PG_URL,
253 object_store,
254 )
255 .await
256 .expect("constructor");
257 (backend, root)
258 }
259
260 #[tokio::test]
261 async fn put_object_bytes_if_absent_stores_and_returns_put() {
262 let (backend, _root) = make_backend().await;
263 let key = make_object_key("test-blob");
264 let data = b"hello object".to_vec();
265 let outcome = backend
266 .put_object_bytes_if_absent(&key, data.clone())
267 .expect("put");
268 assert_eq!(outcome, PutOutcome::Inserted);
269
270 let outcome2 = backend
272 .put_object_bytes_if_absent(&key, data)
273 .expect("put again");
274 assert_eq!(outcome2, PutOutcome::AlreadyExists);
275 }
276
277 #[tokio::test]
278 async fn put_object_bytes_if_absent_empty_bytes() {
279 let (backend, _root) = make_backend().await;
280 let key = make_object_key("empty-blob");
281 let data = Vec::new();
282 let outcome = backend
283 .put_object_bytes_if_absent(&key, data)
284 .expect("put empty");
285 assert_eq!(outcome, PutOutcome::Inserted);
286 }
287
288 #[tokio::test]
289 async fn put_object_bytes_overwrite_stores_without_error() {
290 let (backend, _root) = make_backend().await;
291 let key = make_object_key("overwrite-blob");
292 let data1 = b"first version".to_vec();
293 let data2 = b"second version".to_vec();
294
295 backend
296 .put_object_bytes_overwrite(&key, data1)
297 .expect("first write");
298 backend
299 .put_object_bytes_overwrite(&key, data2)
300 .expect("overwrite");
301 }
303
304 #[tokio::test]
305 async fn copy_object_if_absent_copies_key() {
306 let (backend, _root) = make_backend().await;
307 let src = make_object_key("source-blob");
308 let dst = make_object_key("dest-blob");
309 let data = b"copy me".to_vec();
310
311 backend
312 .put_object_bytes_if_absent(&src, data)
313 .expect("put source");
314 let outcome = backend.copy_object_if_absent(&src, &dst).expect("copy");
315 assert_eq!(outcome, PutOutcome::Inserted);
316
317 let outcome2 = backend
319 .copy_object_if_absent(&src, &dst)
320 .expect("copy again");
321 assert_eq!(outcome2, PutOutcome::AlreadyExists);
322 }
323
324 #[tokio::test]
325 async fn put_sha256_addressed_object_bytes_if_absent_stores_at_canonical_key() {
326 let (backend, _root) = make_backend().await;
327 let data = b"sha256 addressed content".to_vec();
328 let digest_hex = hex::encode(sha2::Sha256::digest(&data));
329 let canonical_key = shared_sha256_object_key(&digest_hex).unwrap();
330 let user_key = make_object_key("user-named-blob");
331
332 let outcome = backend
334 .put_sha256_addressed_object_bytes_if_absent(&canonical_key, &digest_hex, data.clone())
335 .expect("put");
336 assert_eq!(outcome, PutOutcome::Inserted);
337
338 let outcome2 = backend
343 .put_sha256_addressed_object_bytes_if_absent(&user_key, &digest_hex, data)
344 .expect("put with user key");
345 assert!(
348 outcome2 == PutOutcome::Inserted || outcome2 == PutOutcome::AlreadyExists,
349 "expected Inserted or AlreadyExists, got {outcome2:?}"
350 );
351 }
352
353 #[tokio::test]
354 async fn put_sha256_addressed_object_file_stores_from_path() {
355 let (backend, _root) = make_backend().await;
356 let data = b"file content for sha256 addressing";
357 let digest_hex = hex::encode(sha2::Sha256::digest(data));
358 let canonical_key = shared_sha256_object_key(&digest_hex).unwrap();
359 let integrity = ObjectIntegrity::new(
360 shardline_protocol::ShardlineHash::from_bytes(*blake3::hash(data).as_bytes()),
361 data.len() as u64,
362 );
363
364 let tmpfile = tempfile::NamedTempFile::new().expect("temp file");
366 std::fs::write(tmpfile.path(), data).expect("write temp file");
367
368 let outcome = backend
369 .put_sha256_addressed_object_file(
370 &canonical_key,
371 &digest_hex,
372 tmpfile.path(),
373 &integrity,
374 )
375 .expect("put from file");
376 assert_eq!(outcome, PutOutcome::Inserted);
377 }
378
379 #[tokio::test]
380 async fn put_sha256_addressed_object_file_with_matching_digest_succeeds() {
381 let (backend, _root) = make_backend().await;
382 let data = b"file content for sha256 addressing";
383 let digest_hex = hex::encode(sha2::Sha256::digest(data));
384 let canonical_key = shared_sha256_object_key(&digest_hex).unwrap();
385 let integrity = ObjectIntegrity::new(
386 shardline_protocol::ShardlineHash::from_bytes(*blake3::hash(data).as_bytes()),
387 data.len() as u64,
388 );
389
390 let tmpfile = tempfile::NamedTempFile::new().expect("temp file");
391 std::fs::write(tmpfile.path(), data).expect("write temp file");
392
393 let outcome = backend
394 .put_sha256_addressed_object_file(
395 &canonical_key,
396 &digest_hex,
397 tmpfile.path(),
398 &integrity,
399 )
400 .expect("put from file with matching digest");
401 assert_eq!(outcome, PutOutcome::Inserted);
402 }
403
404 #[tokio::test]
405 async fn upload_xorb_rejects_invalid_body() {
406 let (backend, _root) = make_backend().await;
407 let data = b"not a valid xorb body";
409 let hash = hex::encode(blake3::hash(data).as_bytes());
410 let result = backend.upload_xorb(&hash, Bytes::from(data.to_vec())).await;
411 assert!(
412 result.is_err(),
413 "upload_xorb should reject invalid xorb data"
414 );
415 }
416
417 #[tokio::test]
418 async fn upload_xorb_rejects_hash_mismatch() {
419 let (backend, _root) = make_backend().await;
420 let data = b"some content";
421 let wrong_hash = "ab".repeat(32);
423 let result = backend
424 .upload_xorb(&wrong_hash, Bytes::from(data.to_vec()))
425 .await;
426 assert!(result.is_err(), "upload_xorb should reject hash mismatch");
427 }
428
429 #[tokio::test]
430 async fn upload_xorb_rejects_empty_body() {
431 let (backend, _root) = make_backend().await;
432 let data = b"";
433 let hash = hex::encode(blake3::hash(data).as_bytes());
434 let result = backend.upload_xorb(&hash, Bytes::from(data.to_vec())).await;
435 assert!(result.is_err(), "upload_xorb should reject empty body");
436 }
437
438 #[tokio::test]
439 async fn upload_xorb_accepts_valid_xorb() {
440 let (backend, _root) = make_backend().await;
441 let content = b"hello xorb world";
443 let (xorb_bytes, expected_hash) = single_chunk_xorb(content);
444 let result = backend.upload_xorb(&expected_hash, xorb_bytes).await;
445 assert!(result.is_ok(), "upload_xorb should accept a valid xorb");
446 }
447
448 #[tokio::test]
449 async fn upload_xorb_rejects_valid_xorb_with_wrong_hash() {
450 let (backend, _root) = make_backend().await;
451 let content = b"valid content but wrong hash declared";
452 let (xorb_bytes, _actual_hash) = single_chunk_xorb(content);
453 let wrong_hash = "ff".repeat(32);
454 let result = backend.upload_xorb(&wrong_hash, xorb_bytes).await;
455 assert!(
456 result.is_err(),
457 "upload_xorb should reject hash mismatch even for valid xorb"
458 );
459 }
460
461 #[tokio::test]
462 async fn put_object_bytes_overwrite_non_existent_then_existing() {
463 let (backend, _root) = make_backend().await;
464 let key = make_object_key("overwrite-new");
465 let data = b"first data".to_vec();
467 backend
468 .put_object_bytes_overwrite(&key, data)
469 .expect("overwrite non-existent");
470
471 let data2 = b"replacement data".to_vec();
473 backend
474 .put_object_bytes_overwrite(&key, data2)
475 .expect("overwrite existing");
476
477 let meta = backend
479 .object_store()
480 .metadata(&key)
481 .expect("metadata after overwrite");
482 assert!(meta.is_some());
483 }
484
485 #[tokio::test]
486 async fn put_sha256_addressed_object_file_with_user_key() {
487 let (backend, _root) = make_backend().await;
488 let data = b"content for user-key test";
489 let digest_hex = hex::encode(sha2::Sha256::digest(data));
490 let canonical_key = shared_sha256_object_key(&digest_hex).unwrap();
491 let user_key = make_object_key("user-named-file-key");
493 let integrity = ObjectIntegrity::new(
494 shardline_protocol::ShardlineHash::from_bytes(*blake3::hash(data).as_bytes()),
495 data.len() as u64,
496 );
497
498 let tmpfile = tempfile::NamedTempFile::new().expect("temp file");
500 std::fs::write(tmpfile.path(), data).expect("write temp file");
501
502 let outcome = backend
504 .put_sha256_addressed_object_file(&user_key, &digest_hex, tmpfile.path(), &integrity)
505 .expect("put with user key");
506 assert_eq!(outcome, PutOutcome::Inserted);
507
508 let meta = backend
511 .object_store()
512 .metadata(&canonical_key)
513 .expect("meta");
514 assert!(
515 meta.is_some(),
516 "canonical key should exist after user-key put"
517 );
518 }
519
520 #[tokio::test]
521 async fn put_object_bytes_if_absent_large_data() {
522 let (backend, _root) = make_backend().await;
523 let key = make_object_key("large-blob");
524 let data = vec![0xABu8; 1_000_000]; let outcome = backend
526 .put_object_bytes_if_absent(&key, data)
527 .expect("put large");
528 assert_eq!(outcome, PutOutcome::Inserted);
529 }
530
531 #[tokio::test]
532 async fn copy_object_if_absent_nonexistent_source_fails() {
533 let (backend, _root) = make_backend().await;
534 let src = make_object_key("non-existent-source");
535 let dst = make_object_key("dest");
536 let result = backend.copy_object_if_absent(&src, &dst);
537 assert!(result.is_err(), "copy of non-existent source should fail");
538 }
539
540 #[tokio::test]
541 async fn put_object_bytes_if_absent_special_key_characters() {
542 let (backend, _root) = make_backend().await;
543 let key = ObjectKey::parse("test/special-v1.0_data.bin").unwrap();
545 let data = b"special key test".to_vec();
546 let outcome = backend
547 .put_object_bytes_if_absent(&key, data)
548 .expect("put special key");
549 assert_eq!(outcome, PutOutcome::Inserted);
550 }
551
552 #[tokio::test]
553 async fn put_object_bytes_if_absent_roundtrip_verify() {
554 let (backend, _root) = make_backend().await;
555 let key = make_object_key("roundtrip-verify");
556 let original = b"verify roundtrip content".to_vec();
557 backend
558 .put_object_bytes_if_absent(&key, original.clone())
559 .expect("put");
560
561 let meta = backend.object_store().metadata(&key).expect("meta");
563 assert!(meta.is_some());
564 let Some(meta) = meta else { return };
565
566 use shardline_storage::ObjectStore;
568 let read_bytes =
569 crate::object_store::read_full_object(&backend.object_store(), &key, meta.length())
570 .expect("read full");
571 assert_eq!(read_bytes, original);
572 }
573
574 #[tokio::test]
575 async fn put_object_bytes_overwrite_large_data() {
576 let (backend, _root) = make_backend().await;
577 let key = make_object_key("overwrite-large");
578 let data = vec![0xCDu8; 500_000];
579 backend
580 .put_object_bytes_overwrite(&key, data)
581 .expect("overwrite large");
582 }
583
584 #[tokio::test]
585 async fn put_object_bytes_if_absent_double_put_same_data() {
586 let (backend, _root) = make_backend().await;
587 let key = make_object_key("double-put-same");
588 let data = b"same data".to_vec();
589
590 let first = backend
591 .put_object_bytes_if_absent(&key, data.clone())
592 .expect("first put");
593 assert_eq!(first, PutOutcome::Inserted);
594
595 let second = backend
597 .put_object_bytes_if_absent(&key, data)
598 .expect("second put");
599 assert_eq!(second, PutOutcome::AlreadyExists);
600 }
601
602 #[tokio::test]
605 async fn upload_file_rejects_invalid_file_id() {
606 let (backend, _root) = make_backend().await;
607 let body = axum::body::Bytes::from(b"content".to_vec());
608 let result = backend.upload_file("/absolute/path", body, None).await;
609 assert!(matches!(result, Err(ServerError::InvalidFileId)));
610 }
611
612 #[tokio::test]
613 async fn upload_file_stream_rejects_invalid_file_id() {
614 let (backend, _root) = make_backend().await;
615 let body = RequestBodyReader::from_bytes(axum::body::Bytes::from(b"content".to_vec()));
616 let result = backend
617 .upload_file_stream("../traverse", body, None, None)
618 .await;
619 assert!(matches!(result, Err(ServerError::InvalidFileId)));
620 }
621
622 #[tokio::test]
623 async fn upload_file_stream_rejects_empty_file_id() {
624 let (backend, _root) = make_backend().await;
625 let body = RequestBodyReader::from_bytes(axum::body::Bytes::from(b"content".to_vec()));
626 let result = backend.upload_file_stream("", body, None, None).await;
627 assert!(matches!(result, Err(ServerError::InvalidFileId)));
628 }
629
630 #[tokio::test]
633 async fn upload_shard_stream_rejects_empty_body() {
634 let (backend, _root) = make_backend().await;
635 let body = RequestBodyReader::from_bytes(axum::body::Bytes::new());
636 let limits = ShardMetadataLimits::new(
637 std::num::NonZeroUsize::new(100).unwrap(),
638 std::num::NonZeroUsize::new(100).unwrap(),
639 std::num::NonZeroUsize::new(100).unwrap(),
640 std::num::NonZeroUsize::new(100).unwrap(),
641 );
642 let result = backend.upload_shard_stream(body, None, limits).await;
643 assert!(result.is_err());
645 }
646
647 #[tokio::test]
650 async fn put_sha256_addressed_object_bytes_if_absent_canonical_equals_user_key() {
651 let (backend, _root) = make_backend().await;
652 let data = b"canonical equals user key".to_vec();
653 let digest_hex = hex::encode(sha2::Sha256::digest(&data));
654 let canonical_key = shared_sha256_object_key(&digest_hex).unwrap();
655
656 let outcome = backend
659 .put_sha256_addressed_object_bytes_if_absent(&canonical_key, &digest_hex, data)
660 .expect("put with matching keys");
661 assert_eq!(outcome, PutOutcome::Inserted);
662 }
663
664 #[tokio::test]
667 async fn put_sha256_addressed_object_file_canonical_equals_user_key() {
668 let (backend, _root) = make_backend().await;
669 let data = b"file canonical equals user key";
670 let digest_hex = hex::encode(sha2::Sha256::digest(data));
671 let canonical_key = shared_sha256_object_key(&digest_hex).unwrap();
672 let integrity = ObjectIntegrity::new(
673 shardline_protocol::ShardlineHash::from_bytes(*blake3::hash(data).as_bytes()),
674 data.len() as u64,
675 );
676
677 let tmpfile = tempfile::NamedTempFile::new().expect("temp file");
678 std::fs::write(tmpfile.path(), data).expect("write temp file");
679
680 let outcome = backend
681 .put_sha256_addressed_object_file(
682 &canonical_key,
683 &digest_hex,
684 tmpfile.path(),
685 &integrity,
686 )
687 .expect("put with matching keys");
688 assert_eq!(outcome, PutOutcome::Inserted);
689 }
690}