Skip to main content

shardline_server/postgres_backend/
read.rs

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    /// Verifies that the local object store and required Postgres metadata tables are
34    /// reachable.
35    ///
36    /// # Errors
37    ///
38    /// Returns [`ServerError`] when local chunk storage is unreadable, the Postgres
39    /// pool cannot execute queries, or required metadata tables are missing.
40    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    /// Loads reconstruction metadata for a file.
71    ///
72    /// # Errors
73    ///
74    /// Returns [`ServerError`] when the file identifier is invalid or the record is
75    /// missing or unreadable.
76    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    /// Loads the logical byte length for a file version.
94    ///
95    /// # Errors
96    ///
97    /// Returns [`ServerError`] when the file identifier is invalid or the record is
98    /// missing or unreadable.
99    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    /// Reconstructs a file into a contiguous byte vector.
112    ///
113    /// # Errors
114    ///
115    /// Returns [`ServerError`] when metadata or chunk bytes cannot be read.
116    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    /// Reads a stored chunk by hash.
135    ///
136    /// # Errors
137    ///
138    /// Returns [`ServerError`] when the hash is invalid or the chunk is missing.
139    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    /// Loads the stored byte length for a chunk object.
221    ///
222    /// # Errors
223    ///
224    /// Returns [`ServerError`] when the hash is invalid or the chunk is missing.
225    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    /// Streams a stored xorb byte range by hash.
249    ///
250    /// # Errors
251    ///
252    /// Returns [`ServerError`] when the hash is invalid, the xorb is missing, or the
253    /// requested byte range cannot be served.
254    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    /// Reads a stored chunk only when it is reachable from a concrete file version.
267    ///
268    /// # Errors
269    ///
270    /// Returns [`ServerError`] when the hash, file identifier, or content hash are
271    /// invalid, when the file version is missing, or when the chunk is not referenced
272    /// by that version.
273    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    /// Loads the stored byte length for a serialized xorb object.
291    ///
292    /// # Errors
293    ///
294    /// Returns [`ServerError`] when the hash is invalid or the xorb is missing.
295    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    /// Store arbitrary bytes as a chunk in the object store and return the hash + key.
492    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    /// Store arbitrary bytes under any object key.
507    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    // ===== Pure function tests =====
518
519    #[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        // A record with `repository_scope: None` should never match any scope.
560        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        // max_connections=1 should be accepted (lazy pool, no actual connect).
750        let result = connect_postgres_metadata_pool("postgres://localhost:5432/test", 1);
751        assert!(result.is_ok() || result.is_err());
752        // Drop the pool explicitly within the tokio context.
753        drop(result);
754    }
755
756    // ===== chunk_length / read_chunk integration tests =====
757
758    #[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); // valid hex but no chunk stored
771        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]; // 64 KB chunk
796        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    // ===== object_length / read_object integration tests =====
803
804    #[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]; // 128 KB
845        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    // ===== delete_object_if_present tests =====
852
853    #[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        // Verify the object is gone.
864        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        // First delete succeeds.
887        assert_eq!(
888            backend
889                .delete_object_if_present(&key)
890                .await
891                .expect("first delete"),
892            DeleteOutcome::Deleted
893        );
894        // Second delete returns NotFound.
895        assert_eq!(
896            backend
897                .delete_object_if_present(&key)
898                .await
899                .expect("second delete"),
900            DeleteOutcome::NotFound
901        );
902    }
903
904    // ===== visit_object_prefix tests =====
905
906    #[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    // ===== list_object_flat_namespace_page tests =====
966
967    #[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        // Request 2 items; the local store may or may not support start_after,
990        // but should at least respect the limit.
991        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    // ===== Special characters and edge cases =====
999
1000    #[tokio::test]
1001    async fn object_keys_with_special_characters() {
1002        let (backend, _root) = make_backend().await;
1003        // Object keys with hyphens, dots, underscores, and slashes.
1004        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        // Non-hex characters or wrong length should produce an error.
1015        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    // ===== Record store tests using LocalRecordStore =====
1027    // These test RecordTraversal + RecordMutation traits directly,
1028    // allowing record-set/scan/delete patterns without postgres.
1029
1030    /// Helper to open a LocalRecordStore in a temp dir.
1031    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        // The latest record should now have v2 data (same file_id).
1122        assert_eq!(parsed.total_bytes, 2048);
1123    }
1124
1125    // record_scan — scan with various prefixes / repository scopes
1126
1127    #[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        // Should contain 3 distinct locators.
1157        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        // Both revisions share the same repository scope key.
1211        assert_eq!(locators.len(), 2);
1212    }
1213
1214    // record_delete — delete existing, delete non-existent
1215
1216    #[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        // Delete it.
1228        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        // Rewrite the same record.
1265        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    // list_latest_version_records — version record listing
1275
1276    #[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        // Write version records.
1284        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    // ===== modified_since_epoch read access =====
1332
1333    #[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        // Should be a reasonable recent timestamp (>0).
1348        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    // ===== latest_record_locator vs version_record_locator =====
1363
1364    #[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        // Latest and version locators should be different for the same record.
1373        assert_ne!(
1374            latest, version,
1375            "latest and version locator keys must differ"
1376        );
1377    }
1378
1379    // ===== Edge: empty records, zero-byte records =====
1380
1381    #[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    // ===== repository_metadata_inventory — various record counts =====
1407
1408    #[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        // Write 3 latest records + 2 version records in the same repository.
1415        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        // Count latest records in the repository.
1437        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        // Count version records in the repository.
1444        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    // ===== Invalid identifier / missing record errors =====
1509
1510    #[tokio::test]
1511    async fn reconstruction_rejects_invalid_file_id() {
1512        let (backend, _root) = make_backend().await;
1513        // read_record() internally calls validate_identifier() which should reject "/".
1514        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        // Should be InvalidFileId or InvalidContentHash.
1541        assert!(result.is_err());
1542    }
1543
1544    // ── ready() ──────────────────────────────────────────────
1545
1546    #[tokio::test]
1547    async fn ready_fails_when_postgres_unreachable() {
1548        let (backend, _root) = make_backend().await;
1549        // The local_root check passes (temp dir exists), but the SQL probe
1550        // fails because no real Postgres is available.
1551        let result = backend.ready().await;
1552        assert!(result.is_err());
1553    }
1554
1555    // ── read_object_stream ───────────────────────────────────
1556
1557    #[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    // ── xorb_length ──────────────────────────────────────────
1583
1584    #[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    // ── reconstruction with content_hash (version locator branch) ──
1604
1605    #[tokio::test]
1606    async fn reconstruction_with_invalid_content_hash() {
1607        let (backend, _root) = make_backend().await;
1608        // An invalid content_hash should be caught before any PG call.
1609        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        // Valid file_id + content_hash → read_record uses version_record_locator
1619        // and then fails on PG.
1620        let result = backend
1621            .reconstruction("test-file.bin", Some(&make_hash('a')), None, None)
1622            .await;
1623        // The error should be from PG, not InvalidFileId or InvalidContentHash.
1624        assert!(result.is_err());
1625        assert!(!matches!(result, Err(ServerError::InvalidFileId)));
1626        assert!(!matches!(result, Err(ServerError::InvalidContentHash)));
1627    }
1628
1629    // ── file_total_bytes with/without content_hash ──────────
1630
1631    #[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    // ── download_file ──────────────────────────────────────
1651
1652    #[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    // ── read_chunk_for_file_version ──────────────────────────
1670
1671    #[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    // ── repository_references_xorb ──────────────────────────
1681
1682    #[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        // Without PG the traversal should propagate an error.
1690        assert!(result.is_err());
1691    }
1692}