1use shardline_index::{
2 FileRecord, PostgresMetadataStoreError, RecordStore, RecordTraversal, RepositoryRecordScope,
3};
4use shardline_protocol::{ByteRange, RepositoryScope};
5use shardline_storage::{DeleteOutcome, ObjectKey, ObjectMetadata, ObjectPrefix, ObjectStore};
6use sqlx::query_scalar;
7use tokio::task;
8
9use crate::{
10 ServerError,
11 chunk_store::chunk_object_key,
12 download_stream::{ServerByteStream, object_byte_range_stream, object_byte_stream},
13 error::IndexError,
14 object_store::{read_full_object, reconstruct_file_record_bytes, visit_object_prefix},
15 record_store::parse_stored_file_record_bytes,
16 validation::{ensure_directory, validate_content_hash, validate_identifier},
17 xet_adapter::{
18 FileReconstructionResponse, build_reconstruction_response, resolve_dedupe_shard_object,
19 xorb_object_key,
20 },
21};
22
23const REQUIRED_METADATA_TABLES: [&str; 6] = [
24 "shardline_file_records",
25 "shardline_file_reconstructions",
26 "shardline_stored_objects",
27 "shardline_dedupe_shards",
28 "shardline_quarantine_candidates",
29 "shardline_retention_holds",
30];
31
32impl super::PostgresBackend {
33 pub async fn ready(&self) -> Result<(), ServerError> {
41 let object_store = self.object_store();
42 if let Some(local_root) = object_store.local_root() {
43 ensure_directory(local_root).await?;
44 } else {
45 let probe_key = ObjectKey::parse("health/probe")
46 .map_err(|_error| ServerError::InvalidContentHash)?;
47 let _object_store_reachable = object_store.metadata(&probe_key)?;
48 }
49 let _probe = query_scalar::<_, i32>("SELECT 1")
50 .fetch_one(self.record_store.pool())
51 .await
52 .map_err(PostgresMetadataStoreError::from)?;
53
54 for table_name in REQUIRED_METADATA_TABLES {
55 let registered_name = query_scalar::<_, Option<String>>("SELECT to_regclass($1)::text")
56 .bind(table_name)
57 .fetch_one(self.record_store.pool())
58 .await
59 .map_err(PostgresMetadataStoreError::from)?;
60 if registered_name.is_none() {
61 return Err(ServerError::Index(
62 IndexError::MissingRequiredMetadataTable(table_name.to_owned()),
63 ));
64 }
65 }
66
67 Ok(())
68 }
69
70 pub async fn reconstruction(
77 &self,
78 file_id: &str,
79 content_hash: Option<&str>,
80 requested_range: Option<ByteRange>,
81 repository_scope: Option<&RepositoryScope>,
82 ) -> Result<FileReconstructionResponse, ServerError> {
83 let record = self
84 .read_record(file_id, content_hash, repository_scope)
85 .await?;
86 Ok(build_reconstruction_response(
87 self.public_base_url(),
88 &record,
89 requested_range,
90 )?)
91 }
92
93 pub async fn file_total_bytes(
100 &self,
101 file_id: &str,
102 content_hash: Option<&str>,
103 repository_scope: Option<&RepositoryScope>,
104 ) -> Result<u64, ServerError> {
105 let record = self
106 .read_record(file_id, content_hash, repository_scope)
107 .await?;
108 Ok(record.total_bytes)
109 }
110
111 pub async fn download_file(
117 &self,
118 file_id: &str,
119 content_hash: Option<&str>,
120 repository_scope: Option<&RepositoryScope>,
121 ) -> Result<Vec<u8>, ServerError> {
122 let record = self
123 .read_record(file_id, content_hash, repository_scope)
124 .await?;
125 let object_store = self.object_store();
126 let server_frontends = self.server_frontends.clone();
127 task::spawn_blocking(move || {
128 reconstruct_file_record_bytes(&object_store, &server_frontends, &record)
129 })
130 .await
131 .map_err(ServerError::BlockingTask)?
132 }
133
134 pub async fn read_chunk(&self, hash_hex: &str) -> Result<Vec<u8>, ServerError> {
140 let object_store = self.object_store();
141 let object_key = chunk_object_key(hash_hex)?;
142 let metadata = object_store.metadata(&object_key)?;
143 let Some(metadata) = metadata else {
144 return Err(ServerError::NotFound);
145 };
146
147 task::spawn_blocking(move || {
148 read_full_object(&object_store, &object_key, metadata.length())
149 })
150 .await
151 .map_err(ServerError::BlockingTask)?
152 }
153
154 pub(crate) async fn object_length(&self, object_key: &ObjectKey) -> Result<u64, ServerError> {
155 let metadata = self.object_store().metadata(object_key)?;
156 let Some(metadata) = metadata else {
157 return Err(ServerError::NotFound);
158 };
159 Ok(metadata.length())
160 }
161
162 pub(crate) async fn read_object(&self, object_key: &ObjectKey) -> Result<Vec<u8>, ServerError> {
163 let object_store = self.object_store();
164 let metadata = object_store.metadata(object_key)?;
165 let Some(metadata) = metadata else {
166 return Err(ServerError::NotFound);
167 };
168 let object_key = object_key.clone();
169 task::spawn_blocking(move || {
170 read_full_object(&object_store, &object_key, metadata.length())
171 })
172 .await
173 .map_err(ServerError::BlockingTask)?
174 }
175
176 pub(crate) async fn read_object_stream(
177 &self,
178 object_key: &ObjectKey,
179 total_length: u64,
180 range: Option<ByteRange>,
181 ) -> Result<ServerByteStream, ServerError> {
182 let object_store = self.object_store();
183 if let Some(range) = range {
184 return object_byte_range_stream(object_store, object_key.clone(), total_length, range)
185 .await;
186 }
187
188 object_byte_stream(object_store, object_key.clone(), total_length).await
189 }
190
191 pub(crate) fn visit_object_prefix<Visitor>(
192 &self,
193 prefix: &ObjectPrefix,
194 visitor: Visitor,
195 ) -> Result<(), ServerError>
196 where
197 Visitor: FnMut(ObjectMetadata) -> Result<(), ServerError>,
198 {
199 visit_object_prefix(&self.object_store(), prefix, visitor)
200 }
201
202 pub(crate) fn list_object_flat_namespace_page(
203 &self,
204 prefix: &ObjectPrefix,
205 start_after: Option<&ObjectKey>,
206 limit: usize,
207 ) -> Result<Vec<ObjectMetadata>, ServerError> {
208 Ok(self
209 .object_store()
210 .list_flat_namespace_page(prefix, start_after, limit)?)
211 }
212
213 pub(crate) async fn delete_object_if_present(
214 &self,
215 object_key: &ObjectKey,
216 ) -> Result<DeleteOutcome, ServerError> {
217 Ok(self.object_store().delete_if_present(object_key)?)
218 }
219
220 pub async fn chunk_length(&self, hash_hex: &str) -> Result<u64, ServerError> {
226 let object_store = self.object_store();
227 let object_key = chunk_object_key(hash_hex)?;
228 let metadata = object_store.metadata(&object_key)?;
229 let Some(metadata) = metadata else {
230 return Err(ServerError::NotFound);
231 };
232
233 Ok(metadata.length())
234 }
235
236 pub(crate) async fn read_dedupe_shard_stream(
237 &self,
238 hash_hex: &str,
239 ) -> Result<(ServerByteStream, u64), ServerError> {
240 let object_store = self.object_store();
241 let (object_key, total_length) =
242 resolve_dedupe_shard_object(&self.index_store, &object_store, hash_hex).await?;
243 let byte_stream = object_byte_stream(object_store, object_key, total_length).await?;
244
245 Ok((byte_stream, total_length))
246 }
247
248 pub(crate) async fn read_xorb_range_stream(
255 &self,
256 hash_hex: &str,
257 total_length: u64,
258 range: ByteRange,
259 ) -> Result<ServerByteStream, ServerError> {
260 let object_store = self.object_store();
261 let object_key = xorb_object_key(hash_hex)?;
262
263 object_byte_range_stream(object_store, object_key, total_length, range).await
264 }
265
266 pub async fn read_chunk_for_file_version(
274 &self,
275 hash_hex: &str,
276 file_id: &str,
277 content_hash: &str,
278 repository_scope: Option<&RepositoryScope>,
279 ) -> Result<Vec<u8>, ServerError> {
280 let record = self
281 .read_record(file_id, Some(content_hash), repository_scope)
282 .await?;
283 if !record.chunks.iter().any(|chunk| chunk.hash == hash_hex) {
284 return Err(ServerError::NotFound);
285 }
286
287 self.read_chunk(hash_hex).await
288 }
289
290 pub async fn xorb_length(&self, hash_hex: &str) -> Result<u64, ServerError> {
296 let object_store = self.object_store();
297 let object_key = xorb_object_key(hash_hex)?;
298 let metadata = object_store.metadata(&object_key)?;
299 let Some(metadata) = metadata else {
300 return Err(ServerError::NotFound);
301 };
302
303 Ok(metadata.length())
304 }
305
306 pub(crate) async fn repository_references_xorb(
307 &self,
308 hash_hex: &str,
309 repository_scope: &RepositoryScope,
310 ) -> Result<bool, ServerError> {
311 repository_references_hash_in_scope(&self.record_store, hash_hex, repository_scope).await
312 }
313
314 async fn read_record(
315 &self,
316 file_id: &str,
317 content_hash: Option<&str>,
318 repository_scope: Option<&RepositoryScope>,
319 ) -> Result<FileRecord, ServerError> {
320 validate_identifier(file_id)?;
321 let probe = FileRecord {
322 file_id: file_id.to_owned(),
323 content_hash: content_hash.unwrap_or_default().to_owned(),
324 total_bytes: 0,
325 chunk_size: 0,
326 repository_scope: repository_scope.cloned(),
327 chunks: Vec::new(),
328 };
329 let locator = if let Some(content_hash) = content_hash {
330 validate_content_hash(content_hash)?;
331 self.record_store.version_record_locator(&probe)
332 } else {
333 self.record_store.latest_record_locator(&probe)
334 };
335 let bytes = RecordTraversal::read_record_bytes(&self.record_store, &locator)
336 .await
337 .map_err(map_record_store_error)?;
338 parse_stored_file_record_bytes(&bytes)
339 }
340}
341
342pub(crate) async fn repository_references_hash_in_scope<RecordAdapter>(
343 record_store: &RecordAdapter,
344 hash_hex: &str,
345 repository_scope: &RepositoryScope,
346) -> Result<bool, ServerError>
347where
348 RecordAdapter: RecordStore + Sync,
349 ServerError: From<RecordAdapter::Error>,
350{
351 let repository = RepositoryRecordScope::from_repository_scope(repository_scope);
352 let mut found = false;
353 RecordTraversal::visit_repository_latest_records(record_store, &repository, |entry| {
354 if found {
355 return Ok::<(), ServerError>(());
356 }
357 if stored_record_references_hash(&entry.bytes, hash_hex, repository_scope)? {
358 found = true;
359 }
360 Ok(())
361 })
362 .await?;
363
364 if found {
365 return Ok(true);
366 }
367
368 RecordTraversal::visit_repository_version_records(record_store, &repository, |entry| {
369 if found {
370 return Ok::<(), ServerError>(());
371 }
372 if stored_record_references_hash(&entry.bytes, hash_hex, repository_scope)? {
373 found = true;
374 }
375 Ok(())
376 })
377 .await?;
378
379 Ok(found)
380}
381
382fn stored_record_references_hash(
383 bytes: &[u8],
384 hash_hex: &str,
385 repository_scope: &RepositoryScope,
386) -> Result<bool, ServerError> {
387 let record = parse_stored_file_record_bytes(bytes)?;
388 if record.repository_scope.as_ref() != Some(repository_scope) {
389 return Ok(false);
390 }
391
392 Ok(record.chunks.iter().any(|chunk| chunk.hash == hash_hex))
393}
394
395pub(crate) fn connect_postgres_metadata_pool(
396 index_postgres_url: &str,
397 max_connections: u32,
398) -> Result<sqlx::PgPool, ServerError> {
399 sqlx::postgres::PgPoolOptions::new()
400 .max_connections(max_connections)
401 .connect_lazy(index_postgres_url)
402 .map_err(PostgresMetadataStoreError::from)
403 .map_err(ServerError::from)
404}
405
406fn map_record_store_error(error: PostgresMetadataStoreError) -> ServerError {
407 match error {
408 PostgresMetadataStoreError::RecordNotFound => ServerError::NotFound,
409 PostgresMetadataStoreError::Sqlx(_)
410 | PostgresMetadataStoreError::Json(_)
411 | PostgresMetadataStoreError::HashParse(_)
412 | PostgresMetadataStoreError::ObjectKey(_)
413 | PostgresMetadataStoreError::Range(_)
414 | PostgresMetadataStoreError::RetentionHold(_)
415 | PostgresMetadataStoreError::QuarantineCandidate(_)
416 | PostgresMetadataStoreError::WebhookDelivery(_)
417 | PostgresMetadataStoreError::IntegerOutOfRange(_)
418 | PostgresMetadataStoreError::InvalidRecordKind
419 | PostgresMetadataStoreError::InvalidRepoType(_) => {
420 ServerError::Index(IndexError::PostgresMetadata(error))
421 }
422 }
423}
424
425#[cfg(test)]
426mod tests {
427 use std::num::NonZeroUsize;
428
429 use super::*;
430 use crate::error::IndexError;
431 use crate::object_store::ServerObjectStore;
432 use serde_json::to_vec;
433 use shardline_index::{FileChunkRecord, FileRecord, LocalRecordStore, RecordMutation};
434 use shardline_protocol::{RepositoryProvider, RepositoryScope};
435 use shardline_storage::{ObjectBody, ObjectIntegrity};
436 use tempfile::TempDir;
437
438 const TEST_PG_URL: &str = "postgres://localhost:5432/test";
439
440 fn test_scope() -> RepositoryScope {
441 RepositoryScope::new(
442 RepositoryProvider::GitHub,
443 "test-owner",
444 "test-repo",
445 Some("main"),
446 )
447 .unwrap()
448 }
449
450 fn make_hash(ch: char) -> String {
451 std::iter::repeat_n(ch, 64).collect()
452 }
453
454 fn file_record_json_bytes(scope: &RepositoryScope, chunk_hash: &str) -> Vec<u8> {
455 let record = FileRecord {
456 file_id: "test/file.bin".into(),
457 content_hash: make_hash('f'),
458 total_bytes: 1024,
459 chunk_size: 256,
460 repository_scope: Some(scope.clone()),
461 chunks: vec![FileChunkRecord {
462 hash: chunk_hash.to_owned(),
463 offset: 0,
464 length: 256,
465 range_start: 0,
466 range_end: 1,
467 packed_start: 0,
468 packed_end: 256,
469 }],
470 };
471 to_vec(&record).unwrap()
472 }
473
474 async fn make_backend() -> (super::super::PostgresBackend, TempDir) {
475 let root = TempDir::new().expect("temp dir");
476 let object_store =
477 ServerObjectStore::local(root.path().join("chunks")).expect("local store");
478 let backend = super::super::PostgresBackend::new_with_object_store_and_upload_parallelism(
479 root.path().to_path_buf(),
480 "http://127.0.0.1:8080".to_owned(),
481 NonZeroUsize::new(65536).unwrap(),
482 NonZeroUsize::new(64).unwrap(),
483 TEST_PG_URL,
484 object_store,
485 )
486 .await
487 .expect("constructor");
488 (backend, root)
489 }
490
491 fn store_chunk(object_store: &ServerObjectStore, data: &[u8]) -> (String, ObjectKey) {
493 let hash = blake3::hash(data);
494 let hash_hex = hex::encode(hash.as_bytes());
495 let object_key = chunk_object_key(&hash_hex).unwrap();
496 let integrity = ObjectIntegrity::new(
497 shardline_protocol::ShardlineHash::from_bytes(*hash.as_bytes()),
498 data.len() as u64,
499 );
500 object_store
501 .put_if_absent(&object_key, ObjectBody::from_vec(data.to_vec()), &integrity)
502 .unwrap();
503 (hash_hex, object_key)
504 }
505
506 fn store_object(object_store: &ServerObjectStore, key: &ObjectKey, data: &[u8]) {
508 let integrity = ObjectIntegrity::new(
509 shardline_protocol::ShardlineHash::from_bytes(*blake3::hash(data).as_bytes()),
510 data.len() as u64,
511 );
512 object_store
513 .put_if_absent(key, ObjectBody::from_vec(data.to_vec()), &integrity)
514 .unwrap();
515 }
516
517 #[test]
520 fn stored_record_references_hash_matching() {
521 let scope = test_scope();
522 let hash = make_hash('b');
523 let bytes = file_record_json_bytes(&scope, &hash);
524
525 assert!(stored_record_references_hash(&bytes, &hash, &scope).unwrap());
526 }
527
528 #[test]
529 fn stored_record_references_hash_different() {
530 let scope = test_scope();
531 let stored_hash = make_hash('b');
532 let queried_hash = make_hash('q');
533 let bytes = file_record_json_bytes(&scope, &stored_hash);
534
535 assert!(!stored_record_references_hash(&bytes, &queried_hash, &scope).unwrap());
536 }
537
538 #[test]
539 fn stored_record_references_hash_different_scope() {
540 let record_scope = test_scope();
541 let other_scope =
542 RepositoryScope::new(RepositoryProvider::GitHub, "other", "other-repo", None).unwrap();
543 let hash = make_hash('b');
544 let bytes = file_record_json_bytes(&record_scope, &hash);
545
546 assert!(!stored_record_references_hash(&bytes, &hash, &other_scope).unwrap());
547 }
548
549 #[test]
550 fn stored_record_references_hash_invalid_json() {
551 let scope = test_scope();
552 let bytes = b"!!! not valid json !!!";
553
554 assert!(stored_record_references_hash(bytes, &make_hash('a'), &scope).is_err());
555 }
556
557 #[test]
558 fn stored_record_references_hash_null_scope() {
559 let scope = test_scope();
561 let hash = make_hash('b');
562 let bytes = {
563 let record = FileRecord {
564 file_id: "test/file.bin".into(),
565 content_hash: make_hash('f'),
566 total_bytes: 1024,
567 chunk_size: 256,
568 repository_scope: None,
569 chunks: vec![FileChunkRecord {
570 hash: hash.clone(),
571 offset: 0,
572 length: 256,
573 range_start: 0,
574 range_end: 1,
575 packed_start: 0,
576 packed_end: 256,
577 }],
578 };
579 to_vec(&record).unwrap()
580 };
581 assert!(!stored_record_references_hash(&bytes, &hash, &scope).unwrap());
582 }
583
584 #[test]
585 fn stored_record_references_hash_empty_chunks() {
586 let scope = test_scope();
587 let hash = make_hash('x');
588 let bytes = {
589 let record = FileRecord {
590 file_id: "test/file.bin".into(),
591 content_hash: make_hash('f'),
592 total_bytes: 0,
593 chunk_size: 0,
594 repository_scope: Some(scope.clone()),
595 chunks: Vec::new(),
596 };
597 to_vec(&record).unwrap()
598 };
599 assert!(!stored_record_references_hash(&bytes, &hash, &scope).unwrap());
600 }
601
602 #[test]
603 fn stored_record_references_hash_hash_not_in_chunks() {
604 let scope = test_scope();
605 let stored_hash = make_hash('a');
606 let queried_hash = make_hash('z');
607 let bytes = file_record_json_bytes(&scope, &stored_hash);
608 assert!(!stored_record_references_hash(&bytes, &queried_hash, &scope).unwrap());
609 }
610
611 #[test]
612 fn map_record_store_error_not_found() {
613 let error = PostgresMetadataStoreError::RecordNotFound;
614 let result = map_record_store_error(error);
615 assert!(matches!(result, ServerError::NotFound));
616 }
617
618 #[test]
619 fn map_record_store_error_sqlx() {
620 let error = PostgresMetadataStoreError::Sqlx(Box::new(sqlx::Error::PoolClosed));
621 let result = map_record_store_error(error);
622 assert!(matches!(
623 result,
624 ServerError::Index(IndexError::PostgresMetadata(_))
625 ));
626 }
627
628 #[test]
629 fn map_record_store_error_json() {
630 let error =
631 PostgresMetadataStoreError::Json(serde_json::from_str::<()>("invalid").unwrap_err());
632 let result = map_record_store_error(error);
633 assert!(matches!(
634 result,
635 ServerError::Index(IndexError::PostgresMetadata(_))
636 ));
637 }
638
639 #[test]
640 fn map_record_store_error_hash_parse() {
641 let error = PostgresMetadataStoreError::HashParse(
642 shardline_protocol::HashParseError::InvalidLength,
643 );
644 let result = map_record_store_error(error);
645 assert!(matches!(
646 result,
647 ServerError::Index(IndexError::PostgresMetadata(_))
648 ));
649 }
650
651 #[test]
652 fn map_record_store_error_object_key() {
653 let error = PostgresMetadataStoreError::ObjectKey(shardline_storage::ObjectKeyError::Empty);
654 let result = map_record_store_error(error);
655 assert!(matches!(
656 result,
657 ServerError::Index(IndexError::PostgresMetadata(_))
658 ));
659 }
660
661 #[test]
662 fn map_record_store_error_range() {
663 let error = PostgresMetadataStoreError::Range(shardline_protocol::RangeError::Inverted);
664 let result = map_record_store_error(error);
665 assert!(matches!(
666 result,
667 ServerError::Index(IndexError::PostgresMetadata(_))
668 ));
669 }
670
671 #[test]
672 fn map_record_store_error_retention_hold() {
673 let error = PostgresMetadataStoreError::RetentionHold(
674 shardline_index::RetentionHoldError::EmptyReason,
675 );
676 let result = map_record_store_error(error);
677 assert!(matches!(
678 result,
679 ServerError::Index(IndexError::PostgresMetadata(_))
680 ));
681 }
682
683 #[test]
684 fn map_record_store_error_quarantine_candidate() {
685 let error = PostgresMetadataStoreError::QuarantineCandidate(
686 shardline_index::QuarantineCandidateError::InvertedTimeline,
687 );
688 let result = map_record_store_error(error);
689 assert!(matches!(
690 result,
691 ServerError::Index(IndexError::PostgresMetadata(_))
692 ));
693 }
694
695 #[test]
696 fn map_record_store_error_webhook_delivery() {
697 let error = PostgresMetadataStoreError::WebhookDelivery(
698 shardline_index::WebhookDeliveryError::EmptyRepositoryOwner,
699 );
700 let result = map_record_store_error(error);
701 assert!(matches!(
702 result,
703 ServerError::Index(IndexError::PostgresMetadata(_))
704 ));
705 }
706
707 #[test]
708 fn map_record_store_error_integer_out_of_range() {
709 let error = PostgresMetadataStoreError::IntegerOutOfRange("test".to_owned());
710 let result = map_record_store_error(error);
711 assert!(matches!(
712 result,
713 ServerError::Index(IndexError::PostgresMetadata(_))
714 ));
715 }
716
717 #[test]
718 fn map_record_store_error_invalid_record_kind() {
719 let error = PostgresMetadataStoreError::InvalidRecordKind;
720 let result = map_record_store_error(error);
721 assert!(matches!(
722 result,
723 ServerError::Index(IndexError::PostgresMetadata(_))
724 ));
725 }
726
727 #[test]
728 fn map_record_store_error_invalid_repo_type() {
729 let error = PostgresMetadataStoreError::InvalidRepoType("unknown".to_owned());
730 let result = map_record_store_error(error);
731 assert!(matches!(
732 result,
733 ServerError::Index(IndexError::PostgresMetadata(_))
734 ));
735 }
736
737 #[test]
738 fn connect_postgres_metadata_pool_empty_url() {
739 assert!(connect_postgres_metadata_pool("", 5).is_err());
740 }
741
742 #[test]
743 fn connect_postgres_metadata_pool_invalid_url() {
744 assert!(connect_postgres_metadata_pool("!!not-a-valid-url!!", 5).is_err());
745 }
746
747 #[tokio::test]
748 async fn connect_postgres_metadata_pool_very_small_connections_succeeds() {
749 let result = connect_postgres_metadata_pool("postgres://localhost:5432/test", 1);
751 assert!(result.is_ok() || result.is_err());
752 drop(result);
754 }
755
756 #[tokio::test]
759 async fn chunk_length_returns_stored_chunk_length() {
760 let (backend, _root) = make_backend().await;
761 let data = b"chunk data for length test";
762 let (hash_hex, _key) = store_chunk(&backend.object_store(), data);
763 let length = backend.chunk_length(&hash_hex).await.expect("chunk_length");
764 assert_eq!(length, data.len() as u64);
765 }
766
767 #[tokio::test]
768 async fn chunk_length_not_found_for_missing_hash() {
769 let (backend, _root) = make_backend().await;
770 let hash_hex = "ab".repeat(32); let result = backend.chunk_length(&hash_hex).await;
772 assert!(matches!(result, Err(ServerError::NotFound)));
773 }
774
775 #[tokio::test]
776 async fn read_chunk_returns_stored_bytes() {
777 let (backend, _root) = make_backend().await;
778 let data = b"read-chunk test data";
779 let (hash_hex, _key) = store_chunk(&backend.object_store(), data);
780 let read_data = backend.read_chunk(&hash_hex).await.expect("read_chunk");
781 assert_eq!(read_data, data);
782 }
783
784 #[tokio::test]
785 async fn read_chunk_not_found_for_missing_hash() {
786 let (backend, _root) = make_backend().await;
787 let hash_hex = "ba".repeat(32);
788 let result = backend.read_chunk(&hash_hex).await;
789 assert!(matches!(result, Err(ServerError::NotFound)));
790 }
791
792 #[tokio::test]
793 async fn read_chunk_with_large_data() {
794 let (backend, _root) = make_backend().await;
795 let data = vec![0xABu8; 65536]; let (hash_hex, _key) = store_chunk(&backend.object_store(), &data);
797 let read_data = backend.read_chunk(&hash_hex).await.expect("read_chunk");
798 assert_eq!(read_data.len(), 65536);
799 assert_eq!(read_data, data);
800 }
801
802 #[tokio::test]
805 async fn object_length_returns_stored_length() {
806 let (backend, _root) = make_backend().await;
807 let key = ObjectKey::parse("test/obj-length").unwrap();
808 let data = b"object length test data";
809 store_object(&backend.object_store(), &key, data);
810 let length = backend.object_length(&key).await.expect("object_length");
811 assert_eq!(length, data.len() as u64);
812 }
813
814 #[tokio::test]
815 async fn object_length_not_found_for_missing_key() {
816 let (backend, _root) = make_backend().await;
817 let key = ObjectKey::parse("test/missing-obj").unwrap();
818 let result = backend.object_length(&key).await;
819 assert!(matches!(result, Err(ServerError::NotFound)));
820 }
821
822 #[tokio::test]
823 async fn read_object_returns_stored_bytes() {
824 let (backend, _root) = make_backend().await;
825 let key = ObjectKey::parse("test/read-obj").unwrap();
826 let data = b"read object test content";
827 store_object(&backend.object_store(), &key, data);
828 let read_data = backend.read_object(&key).await.expect("read_object");
829 assert_eq!(read_data, data);
830 }
831
832 #[tokio::test]
833 async fn read_object_not_found_for_missing_key() {
834 let (backend, _root) = make_backend().await;
835 let key = ObjectKey::parse("test/never-stored").unwrap();
836 let result = backend.read_object(&key).await;
837 assert!(matches!(result, Err(ServerError::NotFound)));
838 }
839
840 #[tokio::test]
841 async fn read_object_roundtrip_large_blob() {
842 let (backend, _root) = make_backend().await;
843 let key = ObjectKey::parse("test/large-blob").unwrap();
844 let data = vec![0x42u8; 131072]; store_object(&backend.object_store(), &key, &data);
846 let read_data = backend.read_object(&key).await.expect("read_object large");
847 assert_eq!(read_data.len(), 131072);
848 assert_eq!(read_data, data);
849 }
850
851 #[tokio::test]
854 async fn delete_object_if_present_deletes_existing() {
855 let (backend, _root) = make_backend().await;
856 let key = ObjectKey::parse("test/delete-existing").unwrap();
857 store_object(&backend.object_store(), &key, b"to be deleted");
858 let outcome = backend
859 .delete_object_if_present(&key)
860 .await
861 .expect("delete");
862 assert_eq!(outcome, DeleteOutcome::Deleted);
863 assert!(matches!(
865 backend.object_length(&key).await,
866 Err(ServerError::NotFound)
867 ));
868 }
869
870 #[tokio::test]
871 async fn delete_object_if_present_missing_returns_not_found() {
872 let (backend, _root) = make_backend().await;
873 let key = ObjectKey::parse("test/never-existed").unwrap();
874 let outcome = backend
875 .delete_object_if_present(&key)
876 .await
877 .expect("delete missing");
878 assert_eq!(outcome, DeleteOutcome::NotFound);
879 }
880
881 #[tokio::test]
882 async fn delete_object_if_present_double_delete() {
883 let (backend, _root) = make_backend().await;
884 let key = ObjectKey::parse("test/double-delete").unwrap();
885 store_object(&backend.object_store(), &key, b"delete me twice");
886 assert_eq!(
888 backend
889 .delete_object_if_present(&key)
890 .await
891 .expect("first delete"),
892 DeleteOutcome::Deleted
893 );
894 assert_eq!(
896 backend
897 .delete_object_if_present(&key)
898 .await
899 .expect("second delete"),
900 DeleteOutcome::NotFound
901 );
902 }
903
904 #[tokio::test]
907 async fn visit_object_prefix_returns_matching_objects() {
908 let (backend, _root) = make_backend().await;
909 let keys = [
910 ObjectKey::parse("test/prefix/a").unwrap(),
911 ObjectKey::parse("test/prefix/b").unwrap(),
912 ObjectKey::parse("test/prefix/sub/c").unwrap(),
913 ObjectKey::parse("test/other/x").unwrap(),
914 ];
915 for key in &keys {
916 store_object(&backend.object_store(), key, b"data");
917 }
918
919 let prefix = ObjectPrefix::parse("test/prefix").unwrap();
920 let mut found: Vec<String> = Vec::new();
921 backend
922 .visit_object_prefix(&prefix, |meta| {
923 found.push(meta.key().as_str().to_owned());
924 Ok(())
925 })
926 .expect("visit");
927 assert_eq!(found.len(), 3, "should find 3 objects under test/prefix/");
928 assert!(found.iter().any(|k: &String| k.contains("test/prefix/a")));
929 assert!(found.iter().any(|k: &String| k.contains("test/prefix/b")));
930 assert!(
931 found
932 .iter()
933 .any(|k: &String| k.contains("test/prefix/sub/c"))
934 );
935 }
936
937 #[tokio::test]
938 async fn visit_object_prefix_empty_when_no_match() {
939 let (backend, _root) = make_backend().await;
940 let prefix = ObjectPrefix::parse("no/such/prefix").unwrap();
941 let mut count = 0_u64;
942 backend
943 .visit_object_prefix(&prefix, |_| {
944 count += 1;
945 Ok(())
946 })
947 .expect("visit empty");
948 assert_eq!(count, 0);
949 }
950
951 #[tokio::test]
952 async fn visit_object_prefix_empty_store() {
953 let (backend, _root) = make_backend().await;
954 let prefix = ObjectPrefix::parse("test").unwrap();
955 let mut count = 0_u64;
956 backend
957 .visit_object_prefix(&prefix, |_| {
958 count += 1;
959 Ok(())
960 })
961 .expect("visit empty store");
962 assert_eq!(count, 0);
963 }
964
965 #[tokio::test]
968 async fn list_object_flat_namespace_page_returns_at_most_limit() {
969 let (backend, _root) = make_backend().await;
970 for i in 0..10 {
971 let key = ObjectKey::parse(&format!("test/page/obj{i:03}")).unwrap();
972 store_object(&backend.object_store(), &key, b"page data");
973 }
974 let prefix = ObjectPrefix::parse("test/page").unwrap();
975 let page = backend
976 .list_object_flat_namespace_page(&prefix, None, 4)
977 .expect("list page");
978 assert_eq!(page.len(), 4);
979 }
980
981 #[tokio::test]
982 async fn list_object_flat_namespace_page_with_limit() {
983 let (backend, _root) = make_backend().await;
984 for i in 0..5 {
985 let key = ObjectKey::parse(&format!("test/page2/obj{i:03}")).unwrap();
986 store_object(&backend.object_store(), &key, b"data");
987 }
988 let prefix = ObjectPrefix::parse("test/page2").unwrap();
989 let page = backend
992 .list_object_flat_namespace_page(&prefix, None, 2)
993 .expect("list with limit");
994 assert!(!page.is_empty());
995 assert!(page.len() <= 2);
996 }
997
998 #[tokio::test]
1001 async fn object_keys_with_special_characters() {
1002 let (backend, _root) = make_backend().await;
1003 let key = ObjectKey::parse("test/special_chars/file-v1.0.bin").unwrap();
1005 let data = b"special key data";
1006 store_object(&backend.object_store(), &key, data);
1007 let read_data = backend.read_object(&key).await.expect("read special key");
1008 assert_eq!(read_data, data);
1009 }
1010
1011 #[tokio::test]
1012 async fn chunk_length_rejects_invalid_hash_format() {
1013 let (backend, _root) = make_backend().await;
1014 let result = backend.chunk_length("not-a-valid-hex-string").await;
1016 assert!(result.is_err());
1017 }
1018
1019 #[tokio::test]
1020 async fn read_chunk_rejects_invalid_hash_format() {
1021 let (backend, _root) = make_backend().await;
1022 let result = backend.read_chunk("!!invalid!!").await;
1023 assert!(result.is_err());
1024 }
1025
1026 fn local_record_store() -> (LocalRecordStore, TempDir) {
1032 let storage = TempDir::new().expect("temp dir");
1033 (
1034 LocalRecordStore::open(storage.path().to_path_buf()),
1035 storage,
1036 )
1037 }
1038
1039 fn test_file_record(file_id: &str, content_hash: &str, scope: &RepositoryScope) -> FileRecord {
1040 FileRecord {
1041 file_id: file_id.to_owned(),
1042 content_hash: content_hash.to_owned(),
1043 total_bytes: 1024,
1044 chunk_size: 256,
1045 repository_scope: Some(scope.clone()),
1046 chunks: Vec::new(),
1047 }
1048 }
1049
1050 #[tokio::test]
1051 async fn record_set_write_and_read_latest() {
1052 let (store, _storage) = local_record_store();
1053 let scope = test_scope();
1054 let record = test_file_record("test/record-set.bin", &make_hash('a'), &scope);
1055
1056 RecordMutation::write_latest_record(&store, &record)
1057 .await
1058 .expect("write latest");
1059
1060 let locator = RecordTraversal::latest_record_locator(&store, &record);
1061 let exists = RecordTraversal::record_locator_exists(&store, &locator)
1062 .await
1063 .expect("exists");
1064 assert!(exists);
1065
1066 let bytes = RecordTraversal::read_record_bytes(&store, &locator)
1067 .await
1068 .expect("read");
1069 let parsed = parse_stored_file_record_bytes(&bytes).expect("parse");
1070 assert_eq!(parsed.file_id, record.file_id);
1071 assert_eq!(parsed.content_hash, record.content_hash);
1072 }
1073
1074 #[tokio::test]
1075 async fn record_set_write_and_read_version() {
1076 let (store, _storage) = local_record_store();
1077 let scope = test_scope();
1078 let record = test_file_record("test/version-record.bin", &make_hash('b'), &scope);
1079
1080 RecordMutation::write_version_record(&store, &record)
1081 .await
1082 .expect("write version");
1083
1084 let locator = RecordTraversal::version_record_locator(&store, &record);
1085 let exists = RecordTraversal::record_locator_exists(&store, &locator)
1086 .await
1087 .expect("exists");
1088 assert!(exists);
1089
1090 let bytes = RecordTraversal::read_record_bytes(&store, &locator)
1091 .await
1092 .expect("read version");
1093 let parsed = parse_stored_file_record_bytes(&bytes).expect("parse");
1094 assert_eq!(parsed.file_id, record.file_id);
1095 assert_eq!(parsed.content_hash, record.content_hash);
1096 }
1097
1098 #[tokio::test]
1099 async fn record_set_overwrite_latest() {
1100 let (store, _storage) = local_record_store();
1101 let scope = test_scope();
1102 let record_v1 = test_file_record("test/overwrite.bin", &make_hash('a'), &scope);
1103 let record_v2 = FileRecord {
1104 content_hash: make_hash('b'),
1105 total_bytes: 2048,
1106 ..record_v1.clone()
1107 };
1108
1109 RecordMutation::write_latest_record(&store, &record_v1)
1110 .await
1111 .expect("write v1");
1112 RecordMutation::write_latest_record(&store, &record_v2)
1113 .await
1114 .expect("write v2 (overwrite)");
1115
1116 let locator = RecordTraversal::latest_record_locator(&store, &record_v1);
1117 let bytes = RecordTraversal::read_record_bytes(&store, &locator)
1118 .await
1119 .expect("read after overwrite");
1120 let parsed = parse_stored_file_record_bytes(&bytes).expect("parse");
1121 assert_eq!(parsed.total_bytes, 2048);
1123 }
1124
1125 #[tokio::test]
1128 async fn record_scan_list_all_latest_locators() {
1129 let (store, _storage) = local_record_store();
1130 let scope_a = test_scope();
1131 let scope_b = RepositoryScope::new(
1132 RepositoryProvider::GitHub,
1133 "other-owner",
1134 "other-repo",
1135 Some("main"),
1136 )
1137 .unwrap();
1138 let rec1 = test_file_record("file1.bin", &make_hash('a'), &scope_a);
1139 let rec2 = test_file_record("file2.bin", &make_hash('b'), &scope_a);
1140 let rec3 = test_file_record("file3.bin", &make_hash('c'), &scope_b);
1141
1142 RecordMutation::write_latest_record(&store, &rec1)
1143 .await
1144 .unwrap();
1145 RecordMutation::write_latest_record(&store, &rec2)
1146 .await
1147 .unwrap();
1148 RecordMutation::write_latest_record(&store, &rec3)
1149 .await
1150 .unwrap();
1151
1152 let locators = RecordTraversal::list_latest_record_locators(&store)
1153 .await
1154 .expect("list all");
1155
1156 assert_eq!(locators.len(), 3);
1158 }
1159
1160 #[tokio::test]
1161 async fn record_scan_repository_scope() {
1162 let (store, _storage) = local_record_store();
1163 let scope = test_scope();
1164 let other_scope =
1165 RepositoryScope::new(RepositoryProvider::GitHub, "team", "other", Some("main"))
1166 .unwrap();
1167 let rec_in_scope = test_file_record("in-scope.bin", &make_hash('a'), &scope);
1168 let rec_other = test_file_record("other.bin", &make_hash('b'), &other_scope);
1169
1170 RecordMutation::write_latest_record(&store, &rec_in_scope)
1171 .await
1172 .unwrap();
1173 RecordMutation::write_latest_record(&store, &rec_other)
1174 .await
1175 .unwrap();
1176
1177 let repo_scope = shardline_index::RepositoryRecordScope::from_repository_scope(&scope);
1178 let repo_locators =
1179 RecordTraversal::list_repository_latest_record_locators(&store, &repo_scope)
1180 .await
1181 .expect("list repo scope");
1182 assert_eq!(repo_locators.len(), 1);
1183 }
1184
1185 #[tokio::test]
1186 async fn record_scan_different_revisions_same_repository() {
1187 let (store, _storage) = local_record_store();
1188 let main_scope = test_scope();
1189 let release_scope = RepositoryScope::new(
1190 RepositoryProvider::GitHub,
1191 "test-owner",
1192 "test-repo",
1193 Some("release"),
1194 )
1195 .unwrap();
1196 let rec_main = test_file_record("main-file.bin", &make_hash('a'), &main_scope);
1197 let rec_release = test_file_record("release-file.bin", &make_hash('b'), &release_scope);
1198
1199 RecordMutation::write_latest_record(&store, &rec_main)
1200 .await
1201 .unwrap();
1202 RecordMutation::write_latest_record(&store, &rec_release)
1203 .await
1204 .unwrap();
1205
1206 let repo_scope = shardline_index::RepositoryRecordScope::from_repository_scope(&main_scope);
1207 let locators = RecordTraversal::list_repository_latest_record_locators(&store, &repo_scope)
1208 .await
1209 .expect("list repo (all revisions)");
1210 assert_eq!(locators.len(), 2);
1212 }
1213
1214 #[tokio::test]
1217 async fn record_delete_existing_locator() {
1218 let (store, _storage) = local_record_store();
1219 let scope = test_scope();
1220 let record = test_file_record("test/delete-me.bin", &make_hash('x'), &scope);
1221
1222 RecordMutation::write_latest_record(&store, &record)
1223 .await
1224 .unwrap();
1225 let locator = RecordTraversal::latest_record_locator(&store, &record);
1226
1227 RecordMutation::delete_record_locator(&store, &locator)
1229 .await
1230 .expect("delete");
1231
1232 let exists = RecordTraversal::record_locator_exists(&store, &locator)
1233 .await
1234 .expect("exists after delete");
1235 assert!(!exists);
1236 }
1237
1238 #[tokio::test]
1239 async fn record_delete_nonexistent_locator_fails() {
1240 let (store, _storage) = local_record_store();
1241 let scope = test_scope();
1242 let record = test_file_record("test/never-written.bin", &make_hash('y'), &scope);
1243 let locator = RecordTraversal::latest_record_locator(&store, &record);
1244
1245 let result = RecordMutation::delete_record_locator(&store, &locator).await;
1246 assert!(result.is_err());
1247 }
1248
1249 #[tokio::test]
1250 async fn record_delete_then_rewrite() {
1251 let (store, _storage) = local_record_store();
1252 let scope = test_scope();
1253 let record = test_file_record("test/recreate.bin", &make_hash('z'), &scope);
1254
1255 RecordMutation::write_latest_record(&store, &record)
1256 .await
1257 .unwrap();
1258 let locator = RecordTraversal::latest_record_locator(&store, &record);
1259
1260 RecordMutation::delete_record_locator(&store, &locator)
1261 .await
1262 .expect("delete");
1263
1264 RecordMutation::write_latest_record(&store, &record)
1266 .await
1267 .unwrap();
1268 let exists = RecordTraversal::record_locator_exists(&store, &locator)
1269 .await
1270 .expect("exists after rewrite");
1271 assert!(exists);
1272 }
1273
1274 #[tokio::test]
1277 async fn list_latest_version_records_with_version_records() {
1278 let (store, _storage) = local_record_store();
1279 let scope = test_scope();
1280 let rec1 = test_file_record("v1.bin", &make_hash('a'), &scope);
1281 let rec2 = test_file_record("v2.bin", &make_hash('b'), &scope);
1282
1283 RecordMutation::write_version_record(&store, &rec1)
1285 .await
1286 .unwrap();
1287 RecordMutation::write_version_record(&store, &rec2)
1288 .await
1289 .unwrap();
1290
1291 let v_locators = RecordTraversal::list_version_record_locators(&store)
1292 .await
1293 .expect("list versions");
1294 assert_eq!(v_locators.len(), 2);
1295 }
1296
1297 #[tokio::test]
1298 async fn list_latest_version_records_empty_when_no_versions() {
1299 let (store, _storage) = local_record_store();
1300 let locators = RecordTraversal::list_version_record_locators(&store)
1301 .await
1302 .expect("list versions empty");
1303 assert!(locators.is_empty());
1304 }
1305
1306 #[tokio::test]
1307 async fn list_latest_version_records_repository_scoped() {
1308 let (store, _storage) = local_record_store();
1309 let scope = test_scope();
1310 let other_scope =
1311 RepositoryScope::new(RepositoryProvider::GitHub, "team", "other", Some("main"))
1312 .unwrap();
1313 let rec = test_file_record("file.bin", &make_hash('a'), &scope);
1314 let other_rec = test_file_record("other.bin", &make_hash('b'), &other_scope);
1315
1316 RecordMutation::write_version_record(&store, &rec)
1317 .await
1318 .unwrap();
1319 RecordMutation::write_version_record(&store, &other_rec)
1320 .await
1321 .unwrap();
1322
1323 let repo_scope = shardline_index::RepositoryRecordScope::from_repository_scope(&scope);
1324 let locators =
1325 RecordTraversal::list_repository_version_record_locators(&store, &repo_scope)
1326 .await
1327 .expect("list repo versions");
1328 assert_eq!(locators.len(), 1);
1329 }
1330
1331 #[tokio::test]
1334 async fn modified_since_epoch_returns_timestamp_for_stored_record() {
1335 let (store, _storage) = local_record_store();
1336 let scope = test_scope();
1337 let record = test_file_record("test/timestamp.bin", &make_hash('t'), &scope);
1338
1339 RecordMutation::write_latest_record(&store, &record)
1340 .await
1341 .unwrap();
1342 let locator = RecordTraversal::latest_record_locator(&store, &record);
1343
1344 let duration = RecordTraversal::modified_since_epoch(&store, &locator)
1345 .await
1346 .expect("modified_since_epoch");
1347 assert!(duration.as_secs() > 0);
1349 }
1350
1351 #[tokio::test]
1352 async fn modified_since_epoch_fails_for_missing_locator() {
1353 let (store, _storage) = local_record_store();
1354 let scope = test_scope();
1355 let record = test_file_record("test/no-exist.bin", &make_hash('n'), &scope);
1356 let locator = RecordTraversal::latest_record_locator(&store, &record);
1357
1358 let result = RecordTraversal::modified_since_epoch(&store, &locator).await;
1359 assert!(result.is_err());
1360 }
1361
1362 #[tokio::test]
1365 async fn latest_and_version_locators_differ() {
1366 let (store, _storage) = local_record_store();
1367 let scope = test_scope();
1368 let record = test_file_record("test/locator-diff.bin", &make_hash('d'), &scope);
1369
1370 let latest = RecordTraversal::latest_record_locator(&store, &record);
1371 let version = RecordTraversal::version_record_locator(&store, &record);
1372 assert_ne!(
1374 latest, version,
1375 "latest and version locator keys must differ"
1376 );
1377 }
1378
1379 #[tokio::test]
1382 async fn record_set_empty_chunks_list() {
1383 let (store, _storage) = local_record_store();
1384 let scope = test_scope();
1385 let record = FileRecord {
1386 file_id: "test/empty-chunks.bin".into(),
1387 content_hash: make_hash('e'),
1388 total_bytes: 0,
1389 chunk_size: 0,
1390 repository_scope: Some(scope),
1391 chunks: Vec::new(),
1392 };
1393
1394 RecordMutation::write_latest_record(&store, &record)
1395 .await
1396 .unwrap();
1397 let locator = RecordTraversal::latest_record_locator(&store, &record);
1398 let bytes = RecordTraversal::read_record_bytes(&store, &locator)
1399 .await
1400 .expect("read empty-chunks record");
1401 let parsed = parse_stored_file_record_bytes(&bytes).expect("parse");
1402 assert!(parsed.chunks.is_empty());
1403 assert_eq!(parsed.total_bytes, 0);
1404 }
1405
1406 #[tokio::test]
1409 async fn repository_metadata_inventory_counts() {
1410 let (store, _storage) = local_record_store();
1411 let scope = test_scope();
1412 let repo_scope = shardline_index::RepositoryRecordScope::from_repository_scope(&scope);
1413
1414 for i in 0..3 {
1416 let rec = test_file_record(
1417 &format!("inventory/file{i}.bin"),
1418 &make_hash(char::from(b'a' + i)),
1419 &scope,
1420 );
1421 RecordMutation::write_latest_record(&store, &rec)
1422 .await
1423 .unwrap();
1424 }
1425 for i in 0..2 {
1426 let rec = test_file_record(
1427 &format!("inventory/version{i}.bin"),
1428 &make_hash(char::from(b'x' + i)),
1429 &scope,
1430 );
1431 RecordMutation::write_version_record(&store, &rec)
1432 .await
1433 .unwrap();
1434 }
1435
1436 let latest_locators =
1438 RecordTraversal::list_repository_latest_record_locators(&store, &repo_scope)
1439 .await
1440 .expect("list repo latest");
1441 assert_eq!(latest_locators.len(), 3);
1442
1443 let version_locators =
1445 RecordTraversal::list_repository_version_record_locators(&store, &repo_scope)
1446 .await
1447 .expect("list repo versions");
1448 assert_eq!(version_locators.len(), 2);
1449 }
1450
1451 #[tokio::test]
1452 async fn repository_metadata_inventory_empty_repository() {
1453 let (store, _storage) = local_record_store();
1454 let scope = test_scope();
1455 let repo_scope = shardline_index::RepositoryRecordScope::from_repository_scope(&scope);
1456
1457 let latest = RecordTraversal::list_repository_latest_record_locators(&store, &repo_scope)
1458 .await
1459 .expect("list repo latest empty");
1460 assert!(latest.is_empty());
1461
1462 let versions =
1463 RecordTraversal::list_repository_version_record_locators(&store, &repo_scope)
1464 .await
1465 .expect("list repo versions empty");
1466 assert!(versions.is_empty());
1467 }
1468
1469 #[tokio::test]
1470 async fn repository_metadata_inventory_mixed_repositories() {
1471 let (store, _storage) = local_record_store();
1472 let scope_a = test_scope();
1473 let scope_b = RepositoryScope::new(
1474 RepositoryProvider::GitHub,
1475 "owner-b",
1476 "repo-b",
1477 Some("main"),
1478 )
1479 .unwrap();
1480 let repo_a = shardline_index::RepositoryRecordScope::from_repository_scope(&scope_a);
1481 let repo_b = shardline_index::RepositoryRecordScope::from_repository_scope(&scope_b);
1482
1483 let rec_a = test_file_record("a.bin", &make_hash('1'), &scope_a);
1484 let rec_b = test_file_record("b.bin", &make_hash('2'), &scope_b);
1485 RecordMutation::write_latest_record(&store, &rec_a)
1486 .await
1487 .unwrap();
1488 RecordMutation::write_latest_record(&store, &rec_b)
1489 .await
1490 .unwrap();
1491
1492 assert_eq!(
1493 RecordTraversal::list_repository_latest_record_locators(&store, &repo_a)
1494 .await
1495 .expect("repo_a")
1496 .len(),
1497 1
1498 );
1499 assert_eq!(
1500 RecordTraversal::list_repository_latest_record_locators(&store, &repo_b)
1501 .await
1502 .expect("repo_b")
1503 .len(),
1504 1
1505 );
1506 }
1507
1508 #[tokio::test]
1511 async fn reconstruction_rejects_invalid_file_id() {
1512 let (backend, _root) = make_backend().await;
1513 let result = backend
1515 .reconstruction("/absolute/path", None, None, None)
1516 .await;
1517 assert!(matches!(result, Err(ServerError::InvalidFileId)));
1518 }
1519
1520 #[tokio::test]
1521 async fn file_total_bytes_rejects_invalid_file_id() {
1522 let (backend, _root) = make_backend().await;
1523 let result = backend.file_total_bytes("../traverse", None, None).await;
1524 assert!(matches!(result, Err(ServerError::InvalidFileId)));
1525 }
1526
1527 #[tokio::test]
1528 async fn download_file_rejects_invalid_file_id() {
1529 let (backend, _root) = make_backend().await;
1530 let result = backend.download_file("\0null", None, None).await;
1531 assert!(matches!(result, Err(ServerError::InvalidFileId)));
1532 }
1533
1534 #[tokio::test]
1535 async fn read_chunk_for_file_version_rejects_invalid_hash() {
1536 let (backend, _root) = make_backend().await;
1537 let result = backend
1538 .read_chunk_for_file_version("bad-hash", "file.bin", "bad-content-hash", None)
1539 .await;
1540 assert!(result.is_err());
1542 }
1543
1544 #[tokio::test]
1547 async fn ready_fails_when_postgres_unreachable() {
1548 let (backend, _root) = make_backend().await;
1549 let result = backend.ready().await;
1552 assert!(result.is_err());
1553 }
1554
1555 #[tokio::test]
1558 async fn read_object_stream_without_range_ok() {
1559 let (backend, _root) = make_backend().await;
1560 let key = ObjectKey::parse("test/stream/no-range").unwrap();
1561 let data = b"stream data without range";
1562 store_object(&backend.object_store(), &key, data);
1563 let result = backend
1564 .read_object_stream(&key, data.len() as u64, None)
1565 .await;
1566 assert!(result.is_ok());
1567 }
1568
1569 #[tokio::test]
1570 async fn read_object_stream_with_range_ok() {
1571 let (backend, _root) = make_backend().await;
1572 let key = ObjectKey::parse("test/stream/with-range").unwrap();
1573 let data = b"stream data with range spec";
1574 store_object(&backend.object_store(), &key, data);
1575 let range = ByteRange::new(0, data.len() as u64 - 1).unwrap();
1576 let result = backend
1577 .read_object_stream(&key, data.len() as u64, Some(range))
1578 .await;
1579 assert!(result.is_ok());
1580 }
1581
1582 #[tokio::test]
1585 async fn xorb_length_returns_length_for_stored_xorb() {
1586 let (backend, _root) = make_backend().await;
1587 let hash_hex = "ab".repeat(32);
1588 let key = crate::xet_adapter::xorb_object_key(&hash_hex).unwrap();
1589 let data = b"fake xorb content for length test";
1590 store_object(&backend.object_store(), &key, data);
1591 let length = backend.xorb_length(&hash_hex).await.expect("xorb_length");
1592 assert_eq!(length, data.len() as u64);
1593 }
1594
1595 #[tokio::test]
1596 async fn xorb_length_not_found_for_missing_hash() {
1597 let (backend, _root) = make_backend().await;
1598 let hash_hex = "cd".repeat(32);
1599 let result = backend.xorb_length(&hash_hex).await;
1600 assert!(matches!(result, Err(ServerError::NotFound)));
1601 }
1602
1603 #[tokio::test]
1606 async fn reconstruction_with_invalid_content_hash() {
1607 let (backend, _root) = make_backend().await;
1608 let result = backend
1610 .reconstruction("test-file.bin", Some("not-a-valid-hash"), None, None)
1611 .await;
1612 assert!(matches!(result, Err(ServerError::InvalidContentHash)));
1613 }
1614
1615 #[tokio::test]
1616 async fn reconstruction_with_valid_content_hash_fails_without_postgres() {
1617 let (backend, _root) = make_backend().await;
1618 let result = backend
1621 .reconstruction("test-file.bin", Some(&make_hash('a')), None, None)
1622 .await;
1623 assert!(result.is_err());
1625 assert!(!matches!(result, Err(ServerError::InvalidFileId)));
1626 assert!(!matches!(result, Err(ServerError::InvalidContentHash)));
1627 }
1628
1629 #[tokio::test]
1632 async fn file_total_bytes_with_invalid_content_hash() {
1633 let (backend, _root) = make_backend().await;
1634 let result = backend
1635 .file_total_bytes("test-file.bin", Some("bad"), None)
1636 .await;
1637 assert!(matches!(result, Err(ServerError::InvalidContentHash)));
1638 }
1639
1640 #[tokio::test]
1641 async fn file_total_bytes_with_valid_params_fails_without_postgres() {
1642 let (backend, _root) = make_backend().await;
1643 let result = backend
1644 .file_total_bytes("test-file.bin", Some(&make_hash('a')), None)
1645 .await;
1646 assert!(result.is_err());
1647 assert!(!matches!(result, Err(ServerError::InvalidFileId)));
1648 }
1649
1650 #[tokio::test]
1653 async fn download_file_with_invalid_content_hash() {
1654 let (backend, _root) = make_backend().await;
1655 let result = backend
1656 .download_file("test-file.bin", Some("invalid!"), None)
1657 .await;
1658 assert!(matches!(result, Err(ServerError::InvalidContentHash)));
1659 }
1660
1661 #[tokio::test]
1662 async fn download_file_without_content_hash_fails_without_postgres() {
1663 let (backend, _root) = make_backend().await;
1664 let result = backend.download_file("test-file.bin", None, None).await;
1665 assert!(result.is_err());
1666 assert!(!matches!(result, Err(ServerError::InvalidFileId)));
1667 }
1668
1669 #[tokio::test]
1672 async fn read_chunk_for_file_version_rejects_invalid_file_id() {
1673 let (backend, _root) = make_backend().await;
1674 let result = backend
1675 .read_chunk_for_file_version(&make_hash('a'), "/bad/path", &make_hash('b'), None)
1676 .await;
1677 assert!(matches!(result, Err(ServerError::InvalidFileId)));
1678 }
1679
1680 #[tokio::test]
1683 async fn repository_references_xorb_fails_without_postgres() {
1684 let (backend, _root) = make_backend().await;
1685 let scope = test_scope();
1686 let result = backend
1687 .repository_references_xorb(&make_hash('a'), &scope)
1688 .await;
1689 assert!(result.is_err());
1691 }
1692}