cgeorgiaw HF Staff commited on
Commit
f0190da
·
verified ·
1 Parent(s): 1dadabb

Add SQLite accession index and bounded on-demand bucket retrieval

Browse files
FEATURE_BACKLOG.md CHANGED
@@ -2,7 +2,7 @@
2
 
3
  Track desired features that are not yet implemented. Add new ideas here and update each item's status as work progresses. These entries describe future work; implementation details and priorities remain open unless explicitly agreed.
4
 
5
- Last updated: 2026-09-29.
6
 
7
  ## Overview
8
 
@@ -11,8 +11,8 @@ Last updated: 2026-09-29.
11
  | F001 | GenBank inference coverage map | Planned; not implemented |
12
  | F002 | Taxonomy map and navigation | Planned; not implemented |
13
  | F003 | Comparison with RefSeq annotations | Planned; not implemented |
14
- | F004 | Download annotations by accession, individually or in bulk | Planned; only individual sampled segment downloads exist |
15
- | F005 | Fast accession lookup and retrieval across the full dataset | Planned; only a local sample index exists |
16
 
17
  ## F001 — GenBank inference coverage map
18
 
@@ -74,7 +74,7 @@ Decisions and dependencies for implementation:
74
 
75
  **Goal:** Make it easy for users to download annotations for the assemblies or contigs they are interested in.
76
 
77
- Current baseline: the prototype can download one selected sampled segment as Parquet. Accession-level and bulk downloads are future work.
78
 
79
  Desired behavior:
80
 
@@ -97,7 +97,9 @@ Decisions and dependencies for implementation:
97
 
98
  **Goal:** Keep accession searches and annotation retrieval responsive as the app expands beyond its bundled sample.
99
 
100
- Current baseline: the prototype indexes local sample metadata and reads selected Parquet row groups. It does not search the full bucket.
 
 
101
 
102
  Desired behavior:
103
 
 
2
 
3
  Track desired features that are not yet implemented. Add new ideas here and update each item's status as work progresses. These entries describe future work; implementation details and priorities remain open unless explicitly agreed.
4
 
5
+ Last updated: 2026-09-30.
6
 
7
  ## Overview
8
 
 
11
  | F001 | GenBank inference coverage map | Planned; not implemented |
12
  | F002 | Taxonomy map and navigation | Planned; not implemented |
13
  | F003 | Comparison with RefSeq annotations | Planned; not implemented |
14
+ | F004 | Download annotations by accession, individually or in bulk | Partial; individual indexed segment downloads exist |
15
+ | F005 | Fast accession lookup and retrieval across the full dataset | Prototype implemented for three remote objects; full coverage pending |
16
 
17
  ## F001 — GenBank inference coverage map
18
 
 
74
 
75
  **Goal:** Make it easy for users to download annotations for the assemblies or contigs they are interested in.
76
 
77
+ Current baseline: the prototype can download one selected indexed segment as Parquet, fetched from the bucket on demand. Filenames include assembly accession, record name, and coordinates. Complete accession-level and bulk downloads are future work.
78
 
79
  Desired behavior:
80
 
 
97
 
98
  **Goal:** Keep accession searches and annotation retrieval responsive as the app expands beyond its bundled sample.
99
 
100
+ Current baseline: [the remote prototype](README.md#remote-prototype-scope) uses a read-only SQLite metadata index over three published bucket objects, with 53 segments across three assemblies. It fetches selected row groups on demand, checks source content hashes, uses a bounded cache, and displays lookup/retrieval timings. The app supports both older whole-contig records and segmented records. It does not search the full bucket.
101
+
102
+ Remaining work includes expanding index coverage, incremental refresh/publication, pagination beyond 200 matches, and reducing large-row-group transfer costs. `benchmark_remote.py` records first-load and cached performance for the current subset.
103
 
104
  Desired behavior:
105
 
README.md CHANGED
@@ -8,10 +8,10 @@ sdk_version: 6.6.0
8
  python_version: '3.10'
9
  app_file: app.py
10
  pinned: false
11
- short_description: Explore sampled genome annotations by accession ID
12
  ---
13
 
14
- Search a small, real subset of [HuggingFaceBio/genbank-annotations](https://huggingface.co/buckets/HuggingFaceBio/genbank-annotations) by assembly accession, contig accession, or the exact `source_key`. View strand-specific CDS probabilities, inspect provenance, and download original segment rows as Parquet.
15
 
16
  Future features and open design decisions are tracked in the [feature backlog](FEATURE_BACKLOG.md).
17
 
@@ -24,20 +24,54 @@ pip install -r requirements.txt
24
  python app.py
25
  ```
26
 
27
- Open http://localhost:7860. The bundled sample works offline and requires no token at runtime.
28
 
29
  ## Deploy to a Space
30
 
31
  The preview Space is [HuggingFaceBio/genbank-annotation-explorer](https://huggingface.co/spaces/HuggingFaceBio/genbank-annotation-explorer), with private visibility.
32
 
33
- To deploy or update it, upload this directory, including `data/sample.parquet` and `data/manifest.json`. The README metadata configures the runtime, following the [Gradio Spaces documentation](https://huggingface.co/docs/hub/spaces-sdks-gradio). Keep the Space private for the initial preview: the bundled annotations will be accessible to its viewers.
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
34
 
35
  ## Lookup behavior
36
 
37
  - Exact versioned accessions match only that version; unknown versions do not fall back.
38
- - Unversioned accessions return all matching versions in the sample. Input ignores case and surrounding whitespace.
39
- - Assembly lookups return all sampled contigs/segments for that assembly. This sample is not a complete assembly.
40
- - A miss means “not in this sample,” not “absent from GenBank/the bucket.”
41
  - Split contigs stay as separate segments, preserving their source coordinates and segment counts.
42
 
43
  ## Data and coordinates
@@ -62,7 +96,7 @@ The default source object is about 31 MB; automatic downloads are capped at 100
62
  python prepare_sample.py --local-file /path/to/source.parquet
63
  ```
64
 
65
- Source rows are selected intact, without truncating their probabilities. The app never scans or downloads the full bucket at startup. To support arbitrary accessions across the full bucket later, build a separate metadata index mapping accession → bucket object → row group, then retrieve matching probability data on demand.
66
 
67
  ## Checks
68
 
 
8
  python_version: '3.10'
9
  app_file: app.py
10
  pinned: false
11
+ short_description: Find genome annotations with on-demand bucket reads
12
  ---
13
 
14
+ Search an indexed subset of [HuggingFaceBio/genbank-annotations](https://huggingface.co/buckets/HuggingFaceBio/genbank-annotations) by assembly accession, contig accession, or the exact `source_key`. A local SQLite index locates records; probability data stays in the bucket and loads on demand. View strand-specific CDS probabilities, inspect provenance, and download original segment rows as Parquet.
15
 
16
  Future features and open design decisions are tracked in the [feature backlog](FEATURE_BACKLOG.md).
17
 
 
24
  python app.py
25
  ```
26
 
27
+ Open http://localhost:7860. Remote reads require a Hugging Face login (`hf auth login`) or an `HF_TOKEN` with bucket read access. For the original offline sample, run `GENBANK_DATA_MODE=sample python app.py`; that mode requires no token.
28
 
29
  ## Deploy to a Space
30
 
31
  The preview Space is [HuggingFaceBio/genbank-annotation-explorer](https://huggingface.co/spaces/HuggingFaceBio/genbank-annotation-explorer), with private visibility.
32
 
33
+ To deploy or update it, upload this directory, including `data/catalog.sqlite`. Set `HF_TOKEN` as a Space secret with read access to the bucket. Include `data/sample.parquet` and `data/manifest.json` if offline sample mode is desired. Credentials stay on the server and must never be committed. The README metadata configures the runtime, following the [Gradio Spaces documentation](https://huggingface.co/docs/hub/spaces-sdks-gradio). Keep the Space private for the initial preview.
34
+
35
+ ## Remote prototype scope
36
+
37
+ The current snapshot indexes **53 segments, 3 assemblies, and 195,493,177 bases** from three explicitly selected Parquet objects in fungi, protozoa, and vertebrate_other. The SQLite file is **86,016 bytes**. Building it fetched about **208 KB** of file ranges for metadata, rather than the approximately **402 MB** of source files. This subset does not represent the full bucket or complete assemblies.
38
+
39
+ Examples beyond the original bundled sample:
40
+
41
+ - `JBPJTW010000350.1` — fungi, a complete aligned contig.
42
+ - `CM126017.1` — protozoa, an older whole-contig schema without explicit segment columns.
43
+ - `CM057449.1` — Australian lungfish, the final indexed segment of a larger contig, with coordinates beginning at 1,800,000,000. Only this segment is in the index.
44
+
45
+ The app reports lookup time, retrieval time, bytes fetched, and whether the selected data came from the bucket or cache. Byte counts measure returned file ranges, excluding HTTP overhead and retries. Retrieval time excludes plot aggregation and browser rendering.
46
+
47
+ ## Retrieval and storage
48
+
49
+ - SQLite contains accession aliases, metadata, object paths and content hashes, row groups, and row positions. Probability arrays and sequence are excluded.
50
+ - The Space opens the index read-only. Search uses indexed exact aliases, and results are capped at 200 with the total match count shown. Narrow to a contig accession when necessary; pagination is not yet implemented.
51
+ - Remote reads use `HfFileSystem` and PyArrow, fetching the selected row group with all source columns so downloads preserve the original record.
52
+ - An in-memory LRU cache holds at most 512 MiB of Arrow table buffers. Fills are serialized to avoid duplicate concurrent downloads. Temporary read/plot buffers require additional memory; this is not a total process memory limit.
53
+ - Row groups with more than 512 MiB of uncompressed Parquet data are rejected before decoding. Larger source groups will need smaller chunks or another retrieval path.
54
+ - On a cache miss, the object content hash is checked before and after reading. Changed objects are rejected with an index-refresh message. Cached data remains tied to the indexed snapshot. Objects are not versioned by this app; replacement of a source requires rebuilding the index.
55
+ - Older whole-contig records derive display coordinates `[0, aligned_bp_length)`. Their downloaded Parquet rows retain the original schema.
56
+ - Authentication and retrieval errors are shown separately from accession misses. A failed segment selection clears the previous plot and download.
57
+
58
+ ## Rebuild the remote index and benchmark
59
+
60
+ ```bash
61
+ python build_index.py
62
+ python benchmark_remote.py
63
+ ```
64
+
65
+ The index builder reads only metadata columns from the three default files. Use repeated `--source annotations/...parquet` arguments for another bounded set of up to 10 objects. It checks object hashes and SQLite integrity, then replaces the index atomically. Upload the new index and restart/redeploy the app to activate a new snapshot; indexing is not automatic yet.
66
+
67
+ The benchmark fetches one row group per source, validates its selected probability arrays, and compares first-load and cached retrieval plus plot aggregation. It writes `data/benchmark.json` with workspace measurements; these are not a promise of Space latency. The two large source groups require substantial first-load transfers even for a single selected segment. The full collection is never downloaded or indexed at app startup.
68
 
69
  ## Lookup behavior
70
 
71
  - Exact versioned accessions match only that version; unknown versions do not fall back.
72
+ - Unversioned accessions resolve matching versions in the index. Input ignores case and surrounding whitespace.
73
+ - Assembly lookups return indexed contigs/segments for that assembly, subject to the result limit. Indexed coverage may be incomplete.
74
+ - A miss means “not in this indexed subset,” not “absent from GenBank/the bucket.”
75
  - Split contigs stay as separate segments, preserving their source coordinates and segment counts.
76
 
77
  ## Data and coordinates
 
96
  python prepare_sample.py --local-file /path/to/source.parquet
97
  ```
98
 
99
+ Source rows are selected intact, without truncating their probabilities. This command rebuilds only the optional offline sample, independently of the remote SQLite index.
100
 
101
  ## Checks
102
 
app.py CHANGED
@@ -1,8 +1,10 @@
1
  """A small Hugging Face Space for exploring genome annotations by accession."""
2
  from pathlib import Path
 
3
  import os
4
  import re
5
  import tempfile
 
6
 
7
  # Shared compute hosts can have /tmp/gradio owned by a different user.
8
  os.environ.setdefault("GRADIO_TEMP_DIR", str(Path(tempfile.gettempdir()) / f"genbank-explorer-{os.getuid()}"))
@@ -11,23 +13,29 @@ import gradio as gr
11
  import pyarrow.parquet as pq
12
 
13
  from catalog import Catalog
 
14
 
15
 
16
  def build_app(catalog=None):
17
- catalog = catalog or Catalog()
18
- all_ids = list(range(len(catalog.records)))
19
- first = catalog.records[0]
 
 
 
20
 
21
  def search(accession):
22
- ids = catalog.lookup(accession)
 
 
23
  choices = [(f"{catalog.records[i]['record_name']} · {catalog.records[i]['assembly_accession']} · "
24
  f"[{catalog.records[i]['segment_start_bp']:,}, {catalog.records[i]['segment_end_bp']:,})", str(i)) for i in ids]
25
  if not str(accession or "").strip():
26
  message = "Enter an assembly or contig accession. Try an example below."
27
  elif not ids:
28
- message = "No match in this sample. This does not mean the accession is absent from the full bucket."
29
  else:
30
- message = f"Found **{len(ids)} sampled segment(s)**. Assembly results include only contigs loaded in this prototype."
31
  return message, catalog.table(ids), gr.Dropdown(choices=choices, value=str(ids[0]) if ids else None), None
32
 
33
  def select_segment(index):
@@ -36,31 +44,43 @@ def build_app(catalog=None):
36
  record = catalog.records[int(index)]
37
  start = record["segment_start_bp"]
38
  end = record["segment_end_bp"]
39
- plot, step = catalog.window(index, start, end)
40
- return record, start, end, plot, plot_note(step), None
 
 
 
 
41
 
42
- def plot_note(step):
 
43
  return ("Coordinates are **0-based, end-exclusive**. "
44
  + ("Each point is one base." if step == 1 else f"Each point is the mean of up to **{step:,} bases**; short peaks can be smoothed.")
45
- + " Download the segment for the original per-base probabilities.")
 
 
46
 
47
  def update_window(index, start, end):
48
  if index is None:
49
  raise gr.Error("Look up an accession and choose a segment first.")
50
  try:
51
- plot, step = catalog.window(index, start, end)
 
52
  except (ValueError, TypeError, OverflowError) as exc:
53
  raise gr.Error(str(exc)) from exc
54
- return plot, plot_note(step)
55
 
56
  def export(index):
57
  if index is None:
58
  raise gr.Error("Choose a segment first.")
59
- table = catalog.segment_table(index)
 
 
 
60
  assembly = re.sub(r"[^A-Za-z0-9._-]", "_", table["assembly_accession"][0].as_py())
61
  record = re.sub(r"[^A-Za-z0-9._-]", "_", table["record_name"][0].as_py())
62
- start = table["segment_start_bp"][0].as_py()
63
- end = table["segment_end_bp"][0].as_py()
 
64
  filename = f"{assembly}__{record}__{start}-{end}.parquet"
65
  target = Path(tempfile.mkdtemp(prefix="genbank-export-")) / filename
66
  pq.write_table(table, target, compression="zstd")
@@ -68,15 +88,21 @@ def build_app(catalog=None):
68
 
69
  with gr.Blocks(title="GenBank Annotation Explorer", delete_cache=(3600, 3600)) as demo:
70
  gr.Markdown("# GenBank Annotation Explorer\nExplore predicted coding regions by **assembly or contig accession**.")
71
- gr.Markdown(f"**Prototype sample · {len(catalog.records):,} segments · "
72
  f"{len(catalog.manifest['assemblies']):,} assemblies · {catalog.manifest['bases']:,} bases**\n\n"
73
  "Source: [HuggingFaceBio/genbank-annotations](https://huggingface.co/buckets/HuggingFaceBio/genbank-annotations). "
74
- "These are model-predicted CDS probabilities, not curated gene features. Search covers the loaded sample only.")
 
 
75
  with gr.Row():
76
  accession = gr.Textbox(label="Accession ID", placeholder="Assembly (GCA_…) or contig accession", scale=5)
77
  search_button = gr.Button("Find annotations", variant="primary", scale=1)
78
- gr.Examples(examples=[[first["assembly_accession"]], [first["record_name"]],
79
- [first["record_name"].rsplit(".", 1)[0]]], inputs=accession)
 
 
 
 
80
  status = gr.Markdown("Enter an accession or select an example. IDs are case-insensitive; version suffixes are optional.")
81
  results = gr.Dataframe(value=catalog.table([]), interactive=False, label="Matching segments")
82
  segment = gr.Dropdown(choices=[], label="Segment to explore", interactive=True)
@@ -92,9 +118,11 @@ def build_app(catalog=None):
92
  with gr.Row():
93
  download = gr.Button("Prepare segment download")
94
  file = gr.File(label="Original segment annotations (Parquet)", interactive=False)
95
- with gr.Accordion("Accessions available in this sample", open=False):
 
 
96
  gr.Dataframe(value=catalog.table(all_ids), interactive=False)
97
- gr.JSON(value=catalog.manifest, label="Sample provenance")
98
  outputs = [status, results, segment, file]
99
  for event in (search_button.click, accession.submit):
100
  event(search, accession, outputs).then(select_segment, segment, [metadata, start, end, plot, note, file])
 
1
  """A small Hugging Face Space for exploring genome annotations by accession."""
2
  from pathlib import Path
3
+ from contextlib import closing
4
  import os
5
  import re
6
  import tempfile
7
+ import time
8
 
9
  # Shared compute hosts can have /tmp/gradio owned by a different user.
10
  os.environ.setdefault("GRADIO_TEMP_DIR", str(Path(tempfile.gettempdir()) / f"genbank-explorer-{os.getuid()}"))
 
13
  import pyarrow.parquet as pq
14
 
15
  from catalog import Catalog
16
+ from remote_catalog import RemoteCatalog, RemoteReadError
17
 
18
 
19
  def build_app(catalog=None):
20
+ if catalog is None:
21
+ catalog = Catalog() if os.environ.get("GENBANK_DATA_MODE") == "sample" else RemoteCatalog()
22
+ remote_mode = isinstance(catalog, RemoteCatalog)
23
+ all_ids = catalog.browse_ids()
24
+ first = catalog.records[all_ids[0]]
25
+ scope = "indexed subset" if remote_mode else "sample"
26
 
27
  def search(accession):
28
+ began = time.perf_counter()
29
+ ids, total = catalog.find(accession)
30
+ elapsed = time.perf_counter() - began
31
  choices = [(f"{catalog.records[i]['record_name']} · {catalog.records[i]['assembly_accession']} · "
32
  f"[{catalog.records[i]['segment_start_bp']:,}, {catalog.records[i]['segment_end_bp']:,})", str(i)) for i in ids]
33
  if not str(accession or "").strip():
34
  message = "Enter an assembly or contig accession. Try an example below."
35
  elif not ids:
36
+ message = f"No match in this {scope}. This does not mean the accession is absent from the full bucket."
37
  else:
38
+ message = f"Found **{total:,} indexed segment(s)** in {elapsed * 1000:.1f} ms. Showing {len(ids):,}. Assembly coverage may be partial."
39
  return message, catalog.table(ids), gr.Dropdown(choices=choices, value=str(ids[0]) if ids else None), None
40
 
41
  def select_segment(index):
 
44
  record = catalog.records[int(index)]
45
  start = record["segment_start_bp"]
46
  end = record["segment_end_bp"]
47
+ try:
48
+ table, stats = catalog.fetch(index)
49
+ plot, step = catalog.window(index, start, end, table=table)
50
+ except RemoteReadError as exc:
51
+ return record, start, end, None, str(exc), None
52
+ return record, start, end, plot, plot_note(step, stats), None
53
 
54
+ def plot_note(step, stats):
55
+ origin = "local sample" if stats.get("local") else "cache" if stats["cache_hit"] else "bucket"
56
  return ("Coordinates are **0-based, end-exclusive**. "
57
  + ("Each point is one base." if step == 1 else f"Each point is the mean of up to **{step:,} bases**; short peaks can be smoothed.")
58
+ + " Download the segment for the original per-base probabilities.\n\n"
59
+ + f"Loaded from **{origin}** in **{stats['seconds']:.2f} s**"
60
+ + (f" · {stats['bytes_read'] / 1_000_000:.2f} MB fetched." if origin == "bucket" else "."))
61
 
62
  def update_window(index, start, end):
63
  if index is None:
64
  raise gr.Error("Look up an accession and choose a segment first.")
65
  try:
66
+ table, stats = catalog.fetch(index)
67
+ plot, step = catalog.window(index, start, end, table=table)
68
  except (ValueError, TypeError, OverflowError) as exc:
69
  raise gr.Error(str(exc)) from exc
70
+ return plot, plot_note(step, stats)
71
 
72
  def export(index):
73
  if index is None:
74
  raise gr.Error("Choose a segment first.")
75
+ try:
76
+ table = catalog.segment_table(index)
77
+ except RemoteReadError as exc:
78
+ raise gr.Error(str(exc)) from exc
79
  assembly = re.sub(r"[^A-Za-z0-9._-]", "_", table["assembly_accession"][0].as_py())
80
  record = re.sub(r"[^A-Za-z0-9._-]", "_", table["record_name"][0].as_py())
81
+ metadata = catalog.records[int(index)]
82
+ start = metadata["segment_start_bp"]
83
+ end = metadata["segment_end_bp"]
84
  filename = f"{assembly}__{record}__{start}-{end}.parquet"
85
  target = Path(tempfile.mkdtemp(prefix="genbank-export-")) / filename
86
  pq.write_table(table, target, compression="zstd")
 
88
 
89
  with gr.Blocks(title="GenBank Annotation Explorer", delete_cache=(3600, 3600)) as demo:
90
  gr.Markdown("# GenBank Annotation Explorer\nExplore predicted coding regions by **assembly or contig accession**.")
91
+ gr.Markdown(f"**Prototype {scope} · {len(catalog.records):,} segments · "
92
  f"{len(catalog.manifest['assemblies']):,} assemblies · {catalog.manifest['bases']:,} bases**\n\n"
93
  "Source: [HuggingFaceBio/genbank-annotations](https://huggingface.co/buckets/HuggingFaceBio/genbank-annotations). "
94
+ f"These are model-predicted CDS probabilities, not curated gene features. Search covers the {scope} only. "
95
+ + (f"Annotations load on demand from **{len(catalog.manifest['sources'])} bucket files**. "
96
+ f"Index updated {catalog.manifest['created_at'][:10]}." if remote_mode else ""))
97
  with gr.Row():
98
  accession = gr.Textbox(label="Accession ID", placeholder="Assembly (GCA_…) or contig accession", scale=5)
99
  search_button = gr.Button("Find annotations", variant="primary", scale=1)
100
+ examples = [first["assembly_accession"], first["record_name"]]
101
+ if remote_mode:
102
+ with closing(catalog.connect()) as conn:
103
+ examples += [r[0] for r in conn.execute("SELECT record_name FROM segments WHERE id IN (SELECT min(id) FROM segments GROUP BY object_path)")]
104
+ examples += ["JBPJTW010000350.1"]
105
+ gr.Examples(examples=[[e] for e in dict.fromkeys(examples)], inputs=accession)
106
  status = gr.Markdown("Enter an accession or select an example. IDs are case-insensitive; version suffixes are optional.")
107
  results = gr.Dataframe(value=catalog.table([]), interactive=False, label="Matching segments")
108
  segment = gr.Dropdown(choices=[], label="Segment to explore", interactive=True)
 
118
  with gr.Row():
119
  download = gr.Button("Prepare segment download")
120
  file = gr.File(label="Original segment annotations (Parquet)", interactive=False)
121
+ with gr.Accordion("Browse indexed accessions and coverage", open=False):
122
+ gr.Markdown(f"Showing the first {len(all_ids):,} indexed segments. Search an accession to find other indexed records. "
123
+ "An indexed file does not imply complete coverage of its assembly.")
124
  gr.Dataframe(value=catalog.table(all_ids), interactive=False)
125
+ gr.JSON(value=catalog.manifest, label="Index provenance" if remote_mode else "Sample provenance")
126
  outputs = [status, results, segment, file]
127
  for event in (search_button.click, accession.submit):
128
  event(search, accession, outputs).then(select_segment, segment, [metadata, start, end, plot, note, file])
benchmark_remote.py ADDED
@@ -0,0 +1,54 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ """Measure one first-load and cached retrieval per indexed source file."""
2
+ from contextlib import closing
3
+ from datetime import datetime, timezone
4
+ import json
5
+ import platform
6
+ import time
7
+
8
+ import numpy as np
9
+
10
+ from catalog import ROOT, PROBS
11
+ from remote_catalog import RemoteCatalog
12
+
13
+
14
+ def main():
15
+ catalog = RemoteCatalog()
16
+ with closing(catalog.connect()) as conn:
17
+ ids = [r[0] for r in conn.execute("SELECT min(id) FROM segments GROUP BY object_path ORDER BY object_path")]
18
+ results = []
19
+ for index in ids:
20
+ record = catalog.records[index]
21
+ start = time.perf_counter()
22
+ found = catalog.lookup(record["record_name"])
23
+ lookup_ms = (time.perf_counter() - start) * 1000
24
+ assert index in found
25
+ table, cold = catalog.fetch(index)
26
+ # Confirm that actual returned probability arrays agree with indexed coordinates.
27
+ length = record["segment_end_bp"] - record["segment_start_bp"]
28
+ for column in PROBS:
29
+ values = table[column][0].values.to_numpy()
30
+ assert len(values) == length
31
+ assert np.isfinite(values).all()
32
+ assert ((values >= 0) & (values <= 1)).all()
33
+ cached, warm = catalog.fetch(index)
34
+ assert warm["cache_hit"] and warm["bytes_read"] == 0 and table.equals(cached)
35
+ start = time.perf_counter()
36
+ frame, step = catalog.window(index, table=cached)
37
+ plot_seconds = time.perf_counter() - start
38
+ assert len(frame) <= 2400
39
+ result = {"accession": record["record_name"], "assembly": record["assembly_accession"],
40
+ "organism": record["organism_name"], "bases": length,
41
+ "lookup_ms": lookup_ms, "first_load": cold, "cached_load": warm,
42
+ "plot_seconds": plot_seconds, "plot_bin_bp": step}
43
+ results.append(result)
44
+ print(json.dumps(result), flush=True)
45
+ del table, cached, values
46
+ report = {"created_at": datetime.now(timezone.utc).isoformat(),
47
+ "environment": platform.system() + " / Python " + platform.python_version(),
48
+ "note": "Workspace measurements; not Space latency. Empty app cache before each source's first load; upstream caches uncontrolled. Bytes count returned file ranges, excluding HTTP overhead.",
49
+ "index_bytes": catalog.path.stat().st_size, "results": results}
50
+ (ROOT / "data" / "benchmark.json").write_text(json.dumps(report, indent=2) + "\n")
51
+
52
+
53
+ if __name__ == "__main__":
54
+ main()
build_index.py ADDED
@@ -0,0 +1,110 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ """Build a metadata-only SQLite snapshot from an explicit, bounded source list."""
2
+ import argparse
3
+ from datetime import datetime, timezone
4
+ import json
5
+ import os
6
+ from pathlib import Path
7
+ import sqlite3
8
+ import tempfile
9
+ import time
10
+
11
+ from huggingface_hub import HfApi, HfFileSystem, get_token
12
+ import pyarrow.parquet as pq
13
+
14
+ from catalog import ROOT, PROBS, normalize, unversioned, segment_metadata
15
+ from remote_catalog import BUCKET, MeasuredFile, source_info
16
+
17
+ SOURCES = [
18
+ "annotations/fungi/s0000/s0000.chunk-00510.parquet",
19
+ "annotations/protozoa/s0003/s0003.chunk-00017.parquet",
20
+ "annotations/vertebrate_other/s0001/s0001.chunk-00018.parquet",
21
+ ]
22
+
23
+
24
+ def build(output, sources):
25
+ if not 1 <= len(sources) <= 10 or len(set(sources)) != len(sources):
26
+ raise ValueError("Choose 1–10 distinct published Parquet objects.")
27
+ if any(not s.startswith("annotations/") or not s.endswith(".parquet") for s in sources):
28
+ raise ValueError("Only published annotation Parquet objects can be indexed.")
29
+ api, fs = HfApi(token=get_token()), HfFileSystem(token=get_token())
30
+ output = Path(output)
31
+ output.parent.mkdir(parents=True, exist_ok=True)
32
+ fd, name = tempfile.mkstemp(suffix=".sqlite", dir=output.parent)
33
+ os.close(fd)
34
+ temporary = Path(name)
35
+ conn = sqlite3.connect(temporary)
36
+ try:
37
+ conn.executescript("""
38
+ CREATE TABLE segments (
39
+ id INTEGER PRIMARY KEY, assembly_accession TEXT NOT NULL,
40
+ record_name TEXT NOT NULL, segment_start_bp INTEGER NOT NULL,
41
+ segment_end_bp INTEGER NOT NULL, object_path TEXT NOT NULL,
42
+ object_hash TEXT NOT NULL, row_group INTEGER NOT NULL,
43
+ row_in_group INTEGER NOT NULL, metadata_json TEXT NOT NULL);
44
+ CREATE TABLE aliases (alias TEXT NOT NULL, segment_id INTEGER NOT NULL,
45
+ PRIMARY KEY (alias, segment_id)) WITHOUT ROWID;
46
+ CREATE TABLE metadata (key TEXT PRIMARY KEY, value TEXT NOT NULL);
47
+ """)
48
+ manifest = {"bucket_id": BUCKET, "created_at": datetime.now(timezone.utc).isoformat(),
49
+ "scope": "Explicit subset of published objects; not the full bucket or complete assemblies",
50
+ "sources": [], "rows": 0, "bases": 0, "assemblies": [], "divisions": []}
51
+ assemblies, divisions = set(), set()
52
+ for path in sources:
53
+ source = source_info(api, BUCKET, path)
54
+ if source is None:
55
+ raise ValueError(f"Missing source: {path}")
56
+ started = time.perf_counter()
57
+ with MeasuredFile(fs, f"buckets/{BUCKET}/{path}") as remote:
58
+ pf = pq.ParquetFile(remote)
59
+ columns = [c for c in pf.schema_arrow.names if c not in PROBS and c != "sequence"]
60
+ group_sizes = []
61
+ for group in range(pf.num_row_groups):
62
+ group_sizes.append(pf.metadata.row_group(group).total_byte_size)
63
+ rows = pf.read_row_group(group, columns=columns, use_threads=False).to_pylist()
64
+ for row_in_group, record in enumerate(rows):
65
+ record = segment_metadata(record)
66
+ index = manifest["rows"]
67
+ conn.execute("INSERT INTO segments VALUES (?,?,?,?,?,?,?,?,?,?)", (
68
+ index, record["assembly_accession"], record["record_name"],
69
+ record["segment_start_bp"], record["segment_end_bp"], path,
70
+ source.xet_hash, group, row_in_group, json.dumps(record)))
71
+ aliases = {normalize(record["source_key"])}
72
+ for field in ("assembly_accession", "record_name"):
73
+ exact = normalize(record[field])
74
+ aliases.update([exact, unversioned(exact)])
75
+ conn.executemany("INSERT INTO aliases VALUES (?,?)", [(a, index) for a in aliases if a])
76
+ manifest["rows"] += 1
77
+ manifest["bases"] += record["segment_end_bp"] - record["segment_start_bp"]
78
+ assemblies.add(record["assembly_accession"])
79
+ divisions.add(record["division"])
80
+ item = {"path": path, "xet_hash": source.xet_hash, "object_bytes": source.size,
81
+ "rows": pf.metadata.num_rows, "row_groups": pf.num_row_groups,
82
+ "max_uncompressed_row_group_bytes": max(group_sizes),
83
+ "metadata_bytes_fetched": remote.bytes_read, "range_reads": remote.range_reads,
84
+ "index_seconds": round(time.perf_counter() - started, 3)}
85
+ after = source_info(api, BUCKET, path)
86
+ if after is None or after.xet_hash != source.xet_hash:
87
+ raise ValueError(f"Source changed while indexing: {path}")
88
+ manifest["sources"].append(item)
89
+ print(json.dumps(item), flush=True)
90
+ manifest.update(assemblies=sorted(assemblies), divisions=sorted(divisions))
91
+ conn.execute("INSERT INTO metadata VALUES ('manifest',?)", (json.dumps(manifest),))
92
+ conn.commit()
93
+ if conn.execute("PRAGMA integrity_check").fetchone()[0] != "ok":
94
+ raise ValueError("SQLite integrity check failed")
95
+ conn.close()
96
+ os.replace(temporary, output)
97
+ print(json.dumps({"index_bytes": output.stat().st_size, "rows": manifest["rows"],
98
+ "assemblies": len(assemblies), "bases": manifest["bases"]}), flush=True)
99
+ return manifest
100
+ finally:
101
+ conn.close()
102
+ temporary.unlink(missing_ok=True)
103
+
104
+
105
+ if __name__ == "__main__":
106
+ parser = argparse.ArgumentParser(description=__doc__)
107
+ parser.add_argument("--source", action="append", help="Repeat to select up to 10 published objects")
108
+ parser.add_argument("--output", type=Path, default=ROOT / "data" / "catalog.sqlite")
109
+ args = parser.parse_args()
110
+ build(args.output, args.source or SOURCES)
catalog.py CHANGED
@@ -22,6 +22,17 @@ def unversioned(value):
22
  return re.sub(r"\.\d+$", "", value)
23
 
24
 
 
 
 
 
 
 
 
 
 
 
 
25
  class Catalog:
26
  def __init__(self, directory=ROOT / "data"):
27
  directory = Path(directory)
@@ -54,6 +65,20 @@ class Catalog:
54
  def table(self, ids):
55
  return pd.DataFrame([self.records[i] for i in ids], columns=DISPLAY_COLUMNS)
56
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
57
  def segment_table(self, index):
58
  index = int(index)
59
  if not 0 <= index < len(self.records):
@@ -61,14 +86,15 @@ class Catalog:
61
  group, row = self.locations[index]
62
  return pq.ParquetFile(self.path).read_row_group(group).slice(row, 1)
63
 
64
- def window(self, index, start=None, end=None, max_points=1200):
65
  record = self.records[int(index)]
66
  lo, hi = record["segment_start_bp"], record["segment_end_bp"]
67
  start = lo if start is None else int(start)
68
  end = hi if end is None else int(end)
69
  if not lo <= start < end <= hi:
70
  raise ValueError(f"Enter a range within [{lo:,}, {hi:,}) with start < end.")
71
- table = self.segment_table(index)
 
72
  width = end - start
73
  step = max(1, (width + max_points - 1) // max_points)
74
  positions = np.arange(start, end, step)
 
22
  return re.sub(r"\.\d+$", "", value)
23
 
24
 
25
+ def segment_metadata(record):
26
+ """Older outputs store a whole aligned contig without segment columns."""
27
+ record = dict(record)
28
+ if "segment_start_bp" not in record and "segment_end_bp" not in record:
29
+ record.update(segment_start_bp=0, segment_end_bp=record["aligned_bp_length"],
30
+ segment_bp_length=record["aligned_bp_length"], segment_index=0, segment_count=1)
31
+ if not 0 <= record["segment_start_bp"] < record["segment_end_bp"]:
32
+ raise ValueError("Invalid segment coordinates in annotation metadata.")
33
+ return record
34
+
35
+
36
  class Catalog:
37
  def __init__(self, directory=ROOT / "data"):
38
  directory = Path(directory)
 
65
  def table(self, ids):
66
  return pd.DataFrame([self.records[i] for i in ids], columns=DISPLAY_COLUMNS)
67
 
68
+ def find(self, accession, limit=200):
69
+ ids = self.lookup(accession)
70
+ return ids[:limit], len(ids)
71
+
72
+ def browse_ids(self, limit=200):
73
+ return list(range(min(limit, len(self.records))))
74
+
75
+ def fetch(self, index):
76
+ import time
77
+ started = time.perf_counter()
78
+ table = self.segment_table(index)
79
+ return table, {"cache_hit": False, "bytes_read": 0, "range_reads": 0,
80
+ "seconds": time.perf_counter() - started, "local": True}
81
+
82
  def segment_table(self, index):
83
  index = int(index)
84
  if not 0 <= index < len(self.records):
 
86
  group, row = self.locations[index]
87
  return pq.ParquetFile(self.path).read_row_group(group).slice(row, 1)
88
 
89
+ def window(self, index, start=None, end=None, max_points=1200, table=None):
90
  record = self.records[int(index)]
91
  lo, hi = record["segment_start_bp"], record["segment_end_bp"]
92
  start = lo if start is None else int(start)
93
  end = hi if end is None else int(end)
94
  if not lo <= start < end <= hi:
95
  raise ValueError(f"Enter a range within [{lo:,}, {hi:,}) with start < end.")
96
+ if table is None:
97
+ table = self.segment_table(index)
98
  width = end - start
99
  step = max(1, (width + max_points - 1) // max_points)
100
  positions = np.arange(start, end, step)
data/benchmark.json ADDED
@@ -0,0 +1,71 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ {
2
+ "created_at": "2026-09-30T13:29:32.577183+00:00",
3
+ "environment": "Linux / Python 3.10.16",
4
+ "note": "Workspace measurements; not Space latency. Empty app cache before each source's first load; upstream caches uncontrolled. Bytes count returned file ranges, excluding HTTP overhead.",
5
+ "index_bytes": 86016,
6
+ "results": [
7
+ {
8
+ "accession": "JBPJTW010000303.1",
9
+ "assembly": "GCA_055472895.1",
10
+ "organism": "Phakopsora pachyrhizi",
11
+ "bases": 66123,
12
+ "lookup_ms": 0.5395635962486267,
13
+ "first_load": {
14
+ "cache_hit": false,
15
+ "bytes_read": 9178665,
16
+ "range_reads": 22,
17
+ "seconds": 4.416150398552418
18
+ },
19
+ "cached_load": {
20
+ "cache_hit": true,
21
+ "bytes_read": 0,
22
+ "range_reads": 0,
23
+ "seconds": 0.0001310594379901886
24
+ },
25
+ "plot_seconds": 0.0011919979006052017,
26
+ "plot_bin_bp": 56
27
+ },
28
+ {
29
+ "accession": "CM126017.1",
30
+ "assembly": "GCA_052324655.1",
31
+ "organism": "Trichomonas vaginalis",
32
+ "bases": 35626321,
33
+ "lookup_ms": 0.39654597640037537,
34
+ "first_load": {
35
+ "cache_hit": false,
36
+ "bytes_read": 185297497,
37
+ "range_reads": 17,
38
+ "seconds": 8.835233750753105
39
+ },
40
+ "cached_load": {
41
+ "cache_hit": true,
42
+ "bytes_read": 0,
43
+ "range_reads": 0,
44
+ "seconds": 0.00029289815574884415
45
+ },
46
+ "plot_seconds": 1.9516547014936805,
47
+ "plot_bin_bp": 29689
48
+ },
49
+ {
50
+ "accession": "CM057449.1",
51
+ "assembly": "GCA_016271365.2",
52
+ "organism": "Neoceratodus forsteri",
53
+ "bases": 84713947,
54
+ "lookup_ms": 0.3896690905094147,
55
+ "first_load": {
56
+ "cache_hit": false,
57
+ "bytes_read": 186392742,
58
+ "range_reads": 22,
59
+ "seconds": 7.675006306730211
60
+ },
61
+ "cached_load": {
62
+ "cache_hit": true,
63
+ "bytes_read": 0,
64
+ "range_reads": 0,
65
+ "seconds": 0.00034056883305311203
66
+ },
67
+ "plot_seconds": 2.610140291042626,
68
+ "plot_bin_bp": 70595
69
+ }
70
+ ]
71
+ }
data/catalog.sqlite ADDED
Binary file (86 kB). View file
 
remote_catalog.py ADDED
@@ -0,0 +1,154 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ """SQLite metadata lookup with measured, bounded reads from the annotation bucket."""
2
+ from collections import OrderedDict
3
+ from contextlib import closing
4
+ from functools import lru_cache
5
+ import json
6
+ from pathlib import Path
7
+ import sqlite3
8
+ import threading
9
+ import time
10
+
11
+ from huggingface_hub import HfApi, HfFileSystem, get_token
12
+ from huggingface_hub.hf_file_system import HfFileSystemFile
13
+ import pyarrow.parquet as pq
14
+
15
+ from catalog import Catalog, ROOT, PROBS, normalize, segment_metadata
16
+
17
+ BUCKET = "HuggingFaceBio/genbank-annotations"
18
+
19
+
20
+ def source_info(api, bucket, path):
21
+ return next((e for e in api.list_bucket_tree(bucket, path, recursive=False)
22
+ if e.path == path and e.type == "file"), None)
23
+
24
+
25
+ class MeasuredFile(HfFileSystemFile):
26
+ """Count returned range bytes, excluding HTTP headers and retries."""
27
+ def __init__(self, fs, path):
28
+ self.bytes_read = 0
29
+ self.range_reads = 0
30
+ super().__init__(fs, path, mode="rb", block_size=1024 * 1024, cache_type="none")
31
+
32
+ def _fetch_range(self, start, end):
33
+ data = super()._fetch_range(start, end)
34
+ self.bytes_read += len(data)
35
+ self.range_reads += 1
36
+ return data
37
+
38
+
39
+ class RemoteReadError(ValueError):
40
+ pass
41
+
42
+
43
+ class Records:
44
+ def __init__(self, catalog):
45
+ self.catalog = catalog
46
+
47
+ def __len__(self):
48
+ return self.catalog.manifest["rows"]
49
+
50
+ def __getitem__(self, index):
51
+ return self.catalog.record(int(index))
52
+
53
+
54
+ class RemoteCatalog(Catalog):
55
+ def __init__(self, path=ROOT / "data" / "catalog.sqlite", cache_bytes=512 * 1024**2,
56
+ max_group_bytes=512 * 1024**2, api=None, fs=None):
57
+ self.path = Path(path)
58
+ self.cache_limit = cache_bytes
59
+ self.max_group_bytes = max_group_bytes
60
+ self.api = api or HfApi(token=get_token())
61
+ self.fs = fs or HfFileSystem(token=get_token())
62
+ self.cache = OrderedDict()
63
+ self.cache_bytes = 0
64
+ self.lock = threading.Lock()
65
+ with closing(self.connect()) as conn:
66
+ self.manifest = json.loads(conn.execute("SELECT value FROM metadata WHERE key='manifest'").fetchone()[0])
67
+ self.records = Records(self)
68
+
69
+ def connect(self):
70
+ conn = sqlite3.connect(self.path.resolve().as_uri() + "?mode=ro&immutable=1", uri=True)
71
+ conn.row_factory = sqlite3.Row
72
+ return conn
73
+
74
+ @lru_cache(maxsize=512)
75
+ def entry(self, index):
76
+ with closing(self.connect()) as conn:
77
+ row = conn.execute("SELECT * FROM segments WHERE id=?", (int(index),)).fetchone()
78
+ if row is None:
79
+ raise RemoteReadError("Choose a segment from the current index.")
80
+ return dict(row)
81
+
82
+ def record(self, index):
83
+ return json.loads(self.entry(index)["metadata_json"])
84
+
85
+ def browse_ids(self, limit=200):
86
+ with closing(self.connect()) as conn:
87
+ return [r[0] for r in conn.execute("SELECT id FROM segments ORDER BY id LIMIT ?", (limit,))]
88
+
89
+ def find(self, accession, limit=200):
90
+ key = normalize(accession)
91
+ if not key:
92
+ return [], 0
93
+ # Aliases include exact versions and versionless IDs; never strip a query's version.
94
+ with closing(self.connect()) as conn:
95
+ count = conn.execute("SELECT count(*) FROM aliases WHERE alias=?", (key,)).fetchone()[0]
96
+ ids = [r[0] for r in conn.execute(
97
+ "SELECT s.id FROM aliases a JOIN segments s ON s.id=a.segment_id "
98
+ "WHERE a.alias=? ORDER BY s.assembly_accession,s.record_name,s.segment_start_bp,s.id LIMIT ?",
99
+ (key, limit))]
100
+ return ids, count
101
+
102
+ def lookup(self, accession):
103
+ return self.find(accession)[0]
104
+
105
+ def check_source(self, entry):
106
+ source = source_info(self.api, self.manifest["bucket_id"], entry["object_path"])
107
+ if source is None or source.xet_hash != entry["object_hash"]:
108
+ raise RemoteReadError("The bucket object has changed since indexing. Rebuild the index before retrieving this segment.")
109
+
110
+ def fetch(self, index):
111
+ began = time.perf_counter()
112
+ entry = self.entry(int(index))
113
+ key = (entry["object_path"], entry["object_hash"], entry["row_group"])
114
+ stats = {"cache_hit": False, "bytes_read": 0, "range_reads": 0}
115
+ # Serialize cache fills to bound memory and avoid duplicate remote downloads.
116
+ with self.lock:
117
+ if key in self.cache:
118
+ table = self.cache.pop(key)
119
+ self.cache[key] = table
120
+ stats["cache_hit"] = True
121
+ else:
122
+ try:
123
+ self.check_source(entry)
124
+ path = f"buckets/{self.manifest['bucket_id']}/{entry['object_path']}"
125
+ self.fs.invalidate_cache(path)
126
+ with MeasuredFile(self.fs, path) as remote:
127
+ parquet = pq.ParquetFile(remote)
128
+ group = parquet.metadata.row_group(entry["row_group"])
129
+ if group.total_byte_size > self.max_group_bytes:
130
+ raise RemoteReadError("This row group exceeds the prototype's 512 MiB read limit. It needs a smaller storage chunk.")
131
+ table = parquet.read_row_group(entry["row_group"], use_threads=False)
132
+ stats.update(bytes_read=remote.bytes_read, range_reads=remote.range_reads)
133
+ self.check_source(entry)
134
+ except RemoteReadError:
135
+ raise
136
+ except Exception as exc:
137
+ raise RemoteReadError("Could not retrieve annotations from the bucket. Check bucket access and try again; this is not an accession-not-found result.") from exc
138
+ if table.nbytes <= self.cache_limit:
139
+ while self.cache and self.cache_bytes + table.nbytes > self.cache_limit:
140
+ _, old = self.cache.popitem(last=False)
141
+ self.cache_bytes -= old.nbytes
142
+ self.cache[key] = table
143
+ self.cache_bytes += table.nbytes
144
+ result = table.slice(entry["row_in_group"], 1)
145
+ if result.num_rows != 1:
146
+ raise RemoteReadError("The retrieved segment does not match the index. Rebuild the index.")
147
+ actual = segment_metadata(result.select([c for c in result.column_names if c not in PROBS and c != "sequence"]).to_pylist()[0])
148
+ if any(actual[k] != entry[k] for k in ("record_name", "assembly_accession", "segment_start_bp", "segment_end_bp")):
149
+ raise RemoteReadError("The retrieved segment does not match the index. Rebuild the index.")
150
+ stats["seconds"] = time.perf_counter() - began
151
+ return result, stats
152
+
153
+ def segment_table(self, index):
154
+ return self.fetch(index)[0]
requirements.txt CHANGED
@@ -2,3 +2,5 @@ gradio==6.6.0
2
  numpy>=1.26,<3
3
  pandas>=2.2,<3
4
  pyarrow>=18,<24
 
 
 
2
  numpy>=1.26,<3
3
  pandas>=2.2,<3
4
  pyarrow>=18,<24
5
+ huggingface_hub==1.16.0
6
+ fsspec>=2024.6,<2027
tests/test_remote_catalog.py ADDED
@@ -0,0 +1,133 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import io
2
+ import json
3
+ from pathlib import Path
4
+ import sqlite3
5
+ import tempfile
6
+ from types import SimpleNamespace
7
+ import unittest
8
+ from unittest.mock import Mock, patch
9
+
10
+ import pyarrow as pa
11
+ import pyarrow.parquet as pq
12
+
13
+ from remote_catalog import RemoteCatalog, RemoteReadError
14
+ from catalog import segment_metadata
15
+
16
+
17
+ class CountingBuffer(io.BytesIO):
18
+ def __init__(self, data):
19
+ super().__init__(data)
20
+ self.bytes_read = self.range_reads = 0
21
+
22
+ def read(self, n=-1):
23
+ data = super().read(n)
24
+ self.bytes_read += len(data)
25
+ self.range_reads += 1
26
+ return data
27
+
28
+
29
+ class RemoteCatalogTests(unittest.TestCase):
30
+ def setUp(self):
31
+ self.temp = tempfile.TemporaryDirectory()
32
+ self.path = Path(self.temp.name) / 'catalog.sqlite'
33
+ records = []
34
+ for i in range(3):
35
+ version = 1 if i < 2 else 2
36
+ start = 100 if i == 1 else 0
37
+ records.append(dict(assembly_accession=f'GCA_123.{version}', record_name=f'ABC123.{version}',
38
+ source_key=f'GCA_123.{version}|ABC123.{version}', organism_name='Test',
39
+ division='fungi', segment_index=int(i == 1), segment_count=2 if i < 2 else 1,
40
+ segment_start_bp=start, segment_end_bp=start + 5,
41
+ pred_prob_positive_strand_cds=[0., .2, .4, .6, .8],
42
+ pred_prob_negative_strand_cds=[1., .8, .6, .4, .2]))
43
+ self.table = pa.Table.from_pylist(records)
44
+ sink = pa.BufferOutputStream()
45
+ pq.write_table(self.table, sink, row_group_size=1)
46
+ self.parquet = sink.getvalue().to_pybytes()
47
+ with sqlite3.connect(self.path) as conn:
48
+ conn.executescript('''
49
+ CREATE TABLE segments(id INTEGER PRIMARY KEY, assembly_accession TEXT, record_name TEXT,
50
+ segment_start_bp INTEGER, segment_end_bp INTEGER, object_path TEXT, object_hash TEXT,
51
+ row_group INTEGER, row_in_group INTEGER, metadata_json TEXT);
52
+ CREATE TABLE aliases(alias TEXT, segment_id INTEGER, PRIMARY KEY(alias,segment_id));
53
+ CREATE TABLE metadata(key TEXT PRIMARY KEY,value TEXT);
54
+ ''')
55
+ for i, record in enumerate(records):
56
+ metadata = {k: v for k, v in record.items() if not k.startswith('pred_prob_')}
57
+ conn.execute('INSERT INTO segments VALUES(?,?,?,?,?,?,?,?,?,?)', (
58
+ i, record['assembly_accession'], record['record_name'], record['segment_start_bp'],
59
+ record['segment_end_bp'], 'annotations/test.parquet', 'original', i, 0, json.dumps(metadata)))
60
+ aliases = {record['source_key'], 'GCA_123', 'ABC123', record['assembly_accession'], record['record_name']}
61
+ conn.executemany('INSERT INTO aliases VALUES(?,?)', [(a, i) for a in aliases])
62
+ conn.execute('INSERT INTO metadata VALUES(?,?)', ('manifest', json.dumps({'rows': 3, 'bucket_id': 'test/bucket'})))
63
+ self.api, self.fs = Mock(), Mock()
64
+ self.source = patch('remote_catalog.source_info', return_value=SimpleNamespace(xet_hash='original'))
65
+ self.source_mock = self.source.start()
66
+ self.reader = patch('remote_catalog.MeasuredFile', side_effect=lambda *args: CountingBuffer(self.parquet))
67
+ self.reader_mock = self.reader.start()
68
+ self.catalog = RemoteCatalog(self.path, api=self.api, fs=self.fs)
69
+
70
+ def tearDown(self):
71
+ self.reader.stop()
72
+ self.source.stop()
73
+ self.temp.cleanup()
74
+
75
+ def test_lookup_versions_segments_and_limit(self):
76
+ self.assertEqual(self.catalog.find(' abc123.1 '), ([0, 1], 2))
77
+ self.assertEqual(self.catalog.find('GCA_123', limit=2), ([0, 1], 3))
78
+ self.assertEqual(self.catalog.lookup('ABC123.2'), [2])
79
+ self.assertEqual(self.catalog.lookup('ABC123.9'), [])
80
+ self.assertEqual(self.catalog.lookup("' OR 1=1 --"), [])
81
+ self.assertEqual(self.catalog.lookup('GCA_123.1|ABC123.1'), [0, 1])
82
+
83
+ def test_cached_read_and_absolute_window(self):
84
+ table, cold = self.catalog.fetch(1)
85
+ self.assertFalse(cold['cache_hit'])
86
+ self.assertGreater(cold['bytes_read'], 0)
87
+ cached, warm = self.catalog.fetch(1)
88
+ self.assertTrue(warm['cache_hit'])
89
+ self.assertEqual(warm['bytes_read'], 0)
90
+ self.assertTrue(table.equals(cached))
91
+ self.assertEqual(self.reader_mock.call_count, 1)
92
+ frame, step = self.catalog.window(1, table=table)
93
+ self.assertEqual(frame['Position (bp)'].min(), 100)
94
+ self.assertEqual(frame['Position (bp)'].max(), 104)
95
+ self.assertEqual(self.reader_mock.call_count, 1)
96
+
97
+ def test_cache_budget_and_eviction(self):
98
+ self.catalog.cache_limit = pq.ParquetFile(io.BytesIO(self.parquet)).read_row_group(0).nbytes
99
+ self.catalog.fetch(0)
100
+ self.catalog.fetch(1)
101
+ self.assertLessEqual(self.catalog.cache_bytes, self.catalog.cache_limit)
102
+ self.assertEqual(len(self.catalog.cache), 1)
103
+ _, stats = self.catalog.fetch(0)
104
+ self.assertFalse(stats['cache_hit'])
105
+
106
+ def test_stale_source_rejected_before_and_after_read(self):
107
+ self.source_mock.return_value = SimpleNamespace(xet_hash='changed')
108
+ with self.assertRaisesRegex(RemoteReadError, 'changed'):
109
+ self.catalog.fetch(0)
110
+ self.assertEqual(self.reader_mock.call_count, 0)
111
+ self.source_mock.side_effect = [SimpleNamespace(xet_hash='original'), SimpleNamespace(xet_hash='changed')]
112
+ with self.assertRaisesRegex(RemoteReadError, 'changed'):
113
+ self.catalog.fetch(0)
114
+ self.assertEqual(len(self.catalog.cache), 0)
115
+
116
+ def test_oversized_group_and_invalid_id(self):
117
+ self.catalog.max_group_bytes = 1
118
+ with self.assertRaisesRegex(RemoteReadError, 'read limit'):
119
+ self.catalog.fetch(0)
120
+ with self.assertRaises(RemoteReadError):
121
+ self.catalog.fetch(-1)
122
+
123
+ def test_auth_failure_is_not_a_lookup_miss(self):
124
+ self.source_mock.side_effect = RuntimeError('sensitive upstream detail')
125
+ with self.assertRaisesRegex(RemoteReadError, 'Could not retrieve') as error:
126
+ self.catalog.fetch(0)
127
+ self.assertNotIn('sensitive', str(error.exception))
128
+ self.assertEqual(self.catalog.lookup('ABC123.2'), [2])
129
+
130
+ def test_legacy_contig_coordinates(self):
131
+ record = segment_metadata({'aligned_bp_length': 35000, 'record_name': 'OLD.1'})
132
+ self.assertEqual((record['segment_start_bp'], record['segment_end_bp']), (0, 35000))
133
+ self.assertEqual((record['segment_index'], record['segment_count']), (0, 1))