Results and storage¶
Score-store implementations and result consumers should import from altar.results.
results
¶
Stable score-store, schema, and materialized-result API.
Result-backend authors and consumers should import these names from :mod:altar.results. Backend modules
under altar.plugins.data_lake remain implementation details even when their classes are re-exported here.
BigQueryScoreStore
¶
BigQueryScoreStore(
client: Any,
table_id: str,
*,
identity: VariantIdentity | None = None,
lineage_table_id: str | None = None,
)
Bases: ScoreStore
BigQuery-backed score store. One table (project.dataset.table) keyed on variant_id and model_id,
holding the locus identity columns (VariantIdentity) and one column per plugin score column (added by
ensure_schema), clustered for locus pruning.
client is a google.cloud.bigquery.Client or a compatible fake. table_id is
"project.dataset.table". identity is the locus schema, defaulting to VariantIdentity. The client
is injected so this module needs no google-cloud-bigquery at import.
connect
classmethod
¶
connect(
table_id: str,
*,
credentials: Any = None,
lineage_table_id: str | None = None,
identity: VariantIdentity | None = None,
) -> BigQueryScoreStore
Build a store against a real BigQuery client. Requires the altar[bigquery] extra.
add_scores_frame
¶
add_scores_frame(
frame: Any,
*,
model_id: str,
model_name: str,
plugin_identity: PluginRunIdentity,
columns: list[ScoreColumn],
) -> InsertResult
Ingest scores from a pandas DataFrame the caller already holds.
This is the columnar path for a caller that already has the scores as a frame — for example a scoring
pipeline reading its result TSV in pd.read_csv chunks. It skips add_scores' round-trip of rows to
per-row dicts to from_records. frame carries the decomposed locus columns (VariantIdentity:
chr/pos/ref/alt) plus one column per score in columns. variant_id is derived here from the
locus (chr:pos:ref:alt), and genome (the run's genome_build), model_id, and model_name are
stamped on, so the caller supplies only loci and scores. The locus must already be canonical (chr1,
not 1; upper-case alleles); a non-canonical one raises VariantIdentityError before anything loads.
This method is synchronous and blocking, because the columnar bulk load is a blocking BigQuery job, so
it is not part of the async ScoreStore interface. A caller on an event loop should offload it to a
worker thread (e.g. asyncio.to_thread) so a multi-minute load doesn't stall the loop. The
rows-based add_scores is the interface contract; this is a columnar fast path warehouse-scale callers
can opt into.
stage_scores_frame
¶
stage_scores_frame(
frame: Any,
*,
staging_table: str,
model_id: str,
model_name: str,
plugin_identity: PluginRunIdentity,
columns: list[ScoreColumn],
replace: bool = False,
) -> int
Load one result chunk into a job-specific staging table, not the canonical score table.
The first chunk uses replace=True; later chunks append. The staging table is given a seven-day
expiry immediately after its first load, so an interrupted completion cannot leak scratch tables.
Call :meth:promote_staged_scores after the complete file is staged, then
:meth:drop_staged_scores in a finally block.
promote_staged_scores
async
¶
promote_staged_scores(
staging_table: str,
*,
model_id: str,
plugin_identity: PluginRunIdentity,
columns: list[ScoreColumn],
) -> InsertResult
Append a complete, identity-checked staged result to the canonical table.
BigQuery does not enforce unique keys, so concurrent jobs may append the same score key. Read paths defensively select one physical row per key, and deployments may compact old duplicates periodically. The bounded staging source is deduplicated here as an additional guard against repeated scorer rows.
Promotion independently checks that the staging table contains exactly the expected model and plugin run identity, then refreshes the canonical identity for that model. The insert is also filtered to the expected pair, so rows appended to the staging table between validation and insertion cannot cross the identity boundary.
Deliberately do not MERGE score keys against the canonical cache: it is an unpartitioned,
multi-billion-row table, and a source-correlated join cannot prune its clustered blocks. The one
canonical read is clustered by model_id and exists only to enforce run identity; score-key
promotion still scales with the staged batch rather than the full cache.
drop_staged_scores
async
¶
Best-effort caller-facing primitive for removing a score staging table.
get_unscored_variants
async
¶
get_unscored_variants(
variant_ids: Sequence[str],
model_id: str,
*,
variants_relation: str | None = None,
chromosomes: Sequence[str] | None = None,
) -> list[str]
Return the variant_ids not yet scored by model_id, so a scoring pipeline can skip cache hits.
By default the ids are bound inline as @variant_ids. Pass variants_relation (a backtick-quoted
temp-table ref carrying variant_id) to prune at warehouse scale as materialize does, for a variant
set too large to bind as an array param. On the relation path the full input lives only in the temp
table, so the unscored rows are read straight off a LEFT JOIN … IS NULL dedup and variant_ids may be
empty, rather than diffed against an in-memory list. chromosomes then prunes the scores table's chr
cluster key, as in export_unscored; it is ignored on the inline path, whose small array binds
directly.
Both paths filter on the model's genome build, read from its stored run identity. This store keeps one
scientifically compatible identity per model_id, and that identity includes the genome build, so the
filter only prunes the genome cluster key and never hides a compatible score. A model with no stored
identity has no scores, and the query then runs without a genome filter.
get_missing_scores
async
¶
get_missing_scores(
variant_ids: Sequence[str],
model_id: str,
*,
plugin_identity: PluginRunIdentity,
columns: Sequence[ScoreColumn],
reuse_policy: ScoreReusePolicy = ScoreReusePolicy(),
variants_relation: str | None = None,
chromosomes: Sequence[str] | None = None,
) -> list[str]
Filter only candidates, using explicit input genome and accepted lineages.
variants_relation is the warehouse-scale equivalent of the inline ID
list. For large miss sets, use export_unscored to keep them in BigQuery.
prepare_unscored_table
async
¶
prepare_unscored_table(
*,
dest_table: str,
variants_relation: str,
model_id: str,
plugin_identity: PluginRunIdentity,
columns: Sequence[ScoreColumn],
max_variants_per_batch: int,
reuse_policy: ScoreReusePolicy = ScoreReusePolicy(),
chromosomes: Sequence[str] | None = None,
expiration_timestamp: str = "TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 1 DAY)",
variant_eligibility: VariantEligibility | None = None,
) -> int
Snapshot and densely number all accepted-cache misses in one warehouse query.
Returns the miss count. The destination contains canonical locus fields,
variant_id and zero-based scoring_batch. Export each batch with
export_prepared_batch; those exports read only this bounded job table.
The caller owns the candidate relation, destination and task construction.
Pass the model's manifest.variant_eligibility so ineligible candidates
are left out before the cache anti-join, as prepare_scoring_batches does;
they are neither counted nor staged.
export_prepared_batch
async
¶
export_prepared_batch(
*,
prepared_relation: str,
scoring_batch: int,
destination_uri: str,
) -> int
Export one previously prepared miss batch without rescanning canonical scores.
export_unscored
async
¶
export_unscored(
*,
model_id: str,
destination_uri: str,
variants_relation: str,
chromosomes: Sequence[str] | None = None,
scoring_batch: int | None = None,
scoring_batch_count: int | None = None,
score_cache_cutoff: datetime | None = None,
scoring_batch_size: int | None = None,
plugin_identity: PluginRunIdentity | None = None,
columns: Sequence[ScoreColumn] = (),
reuse_policy: ScoreReusePolicy = ScoreReusePolicy(),
variant_eligibility: VariantEligibility | None = None,
) -> int
Export the loci of variants in variants_relation not yet scored by model_id, and return the count
exported. Writes with EXPORT DATA to destination_uri, an already-resolved gs://…-*.tsv wildcard
prefix, staging each model's unscored-variant file. This is a sibling of export_download and
export_prioritized; it differs only in the predicate (LEFT JOIN scores … WHERE s.variant_id IS NULL)
and a flat headerless locus projection.
variants_relation (a backtick-quoted temp-table ref) is required, and unlike materialize's relation
it must carry the decomposed locus columns (VariantIdentity: chr/pos/ref/alt), because the
export projects them; the caller's temp table must produce exactly those columns. chromosomes prunes
the scores table's chr cluster key; omit it for a correct but unpruned scan. The genome filter is
plugin_identity.genome_build when a run identity is given, otherwise the model's stored genome build.
First a parameterized COUNT(*) runs, returning 0 and skipping the export when nothing is unscored.
Then the EXPORT DATA OPTIONS(… CSV, tab, header=false) AS dedup SELECT runs. EXPORT DATA forbids
query params, so genome, model_id, and chromosomes are inlined into the export SQL (they are
controlled internal values) while the count binds them as params.
Pass the model's manifest.variant_eligibility as variant_eligibility so the candidate relation is
restricted to variants the model can score before any miss is counted, batched, or exported. Without it an
SNV-only model's export includes indels, and run_scoring then rejects their result rows.
Resolving the gs:// URI, creating the temp table, and composing the resulting shards stay the caller's
job; the store's job ends at issuing the export. This does not run on the emulator: the EXPORT DATA GCS
sink hangs the job poll, as in export_download, so the generated SQL is checked with a fake client for
shape plus an opt-in smoke test against real BigQuery.
materialize
async
¶
materialize(
*,
variant_ids: Sequence[str],
models: Sequence[ScoredModel],
annotation_sources: Sequence[AnnotationSource],
variants_relation: str | None = None,
) -> AsyncIterator[MaterializedVariant]
Stream one per-variant summary at a time from a single fused query.
The query is one SELECT that LEFT JOINs the scores table to each annotation source's sql_join
JoinSpec. It computes each model's prioritize_predicate().to_sql() server-side as a per-row
is_prioritized column — the same generated SQL the EXPORT paths filter on — and each prioritizing
source's predicate as a boolean column. SQL decides prioritization; the model_scores entries are then
shaped by each model's to_model_score in Python. A parity test checks the result matches the pure
materialize_variant driver row for row, so a caller can trust the pushed-down query.
Streaming: rows arrive ordered by variant_id, so grouping yields one MaterializedVariant at a time
without buffering the whole result, which can span millions of variants. That order is the
ScoreStore.materialize contract: one row per distinct requested variant, ascending by variant_id.
Scope: one MaterializedVariant per distinct requested variant, as with the portable driver. The fused
query covers every variant in variants_relation when one is given. On the inline path it covers the
variants a model scored plus those a fused prioritizing source has a row for (see build_scored_cte), so
a coding-only variant (e.g. an AlphaMissense missense row with no ChromBPNet score) comes back with empty
model_scores and prioritized driven by the source. The remaining requested variants — ones with only
passive annotations, or with no evidence at all — are built in Python with empty model_scores, their
annotations read and prioritizing sources evaluated over one fetch_annotations, and merged into the
sorted stream in variant_id order.
Sources on this backend join in-SQL (sql_join returns a spec); a source on another backend falls back to
annotate(), plus Python evaluation when it prioritizes. A passive column referenced by a model
predicate must come from a fused source, since it has to be in the query for the predicate to push down.
Annotations: each result's annotations carries the fused sources' columns, which the query already
selects for the predicates, plus the other sources' annotate() values, with the same values and the
same absent-when-empty rule as the portable driver. A declared scalar json column that the query returns
as JSON text (a STRING such as TO_JSON_STRING(...)) is decoded to its JSON value first. The unscored
variants outside the fused query read only the fused passive sources again; a fused prioritizing source
has no row for them, and the other sources' values were already read. A relation-scoped call with a
source that cannot fuse raises ValueError unless the job's ids are also passed in variant_ids, since
annotate() needs the ids.
Scoping: variant_ids are bound inline as @variant_ids. Pass variants_relation (a temp-table ref)
instead to scope a warehouse-scale job whose id array would exceed BigQuery's request limit; the SQL
then prunes via variant_id IN (SELECT variant_id FROM <relation>) and variant_ids may be empty. The
two paths return the same result over the same variant set (checked on the emulator).
export_download
async
¶
export_download(
*,
variant_ids: Sequence[str],
models: Sequence[ScoredModel],
annotation_sources: Sequence[AnnotationSource],
destination_uri: str,
rider: DisplayRider | None = None,
variants_relation: str | None = None,
prioritized_only: bool = False,
header: bool = True,
scored_relation: str | None = None,
anno_relation: str | None = None,
) -> str
Export variant_ids as a flat wide TSV to destination_uri (a gs:// prefix) and return it. This is
the download branch of the two exports. prioritized_only restricts to prioritized variants (the same
rollup as export_prioritized, but in the wide-TSV shape) for the prioritized download. header=False
writes headerless parts so a caller doing a server-side GCS compose can prepend a single header (see
download_header) instead of one per shard.
This is a different artifact from materialize: all variants (not prioritized-only), one row per
variant, with each model's score columns pivoted into positional columns (MAX(IF(model_id=…))) rather
than a nested ARRAY<STRUCT>. It is the human-downloadable TSV, not the DB rows. It shares the same
join and pushed-down is_prioritized base (the scored CTE) with materialize, then aggregates:
ANY_VALUE per annotation column, and prioritized = COALESCE(LOGICAL_OR(is_prioritized), FALSE) OR …
over the fused sources. It runs as one EXPORT DATA OPTIONS(... CSV, tab, GZIP) job. BigQuery cannot
export nested or repeated data as CSV, so each annotation column declared repeated or struct is
written as its TO_JSON_STRING text, with a missing value left empty. Such a column must physically be
an ARRAY or STRUCT, not a pre-serialized string. A scalar json column is exported unchanged. The Parquet
and JSON exports keep those columns' native types.
Fused sources only: a flat cloud EXPORT is single-backend, so unlike materialize there is no
annotate() fallback. A source that cannot sql_join raises ValueError.
rider is an optional config-injected DisplayRider — for example most_active_celltype and
most_active_celltype_logfc display columns, the single model-score row picked per variant by an
ordering over the scored namespace. It applies to this aggregated export and the DB export, never to
materialize.
This does not run on the emulator: EXPORT DATA … OPTIONS(uri='gs://…') hangs the job poll on the
emulator's GCS sink, so the generated SQL is checked with a fake client for shape and the actual GCS
write is an opt-in smoke test against real BigQuery. variants_relation scopes as in materialize (a
temp-table ref for a warehouse-scale job; None for inline @variant_ids).
scored_relation reads a pre-materialized scored table (from materialize_scored_table) instead of
rebuilding and rescanning the fused base CTE. When set, the join and pushed-down is_prioritized are not
recomputed here; the export reads the already-computed table, so several exports off one base pay the
base scan (which can be multi-TB) once instead of once each. variants_relation and variant_ids are
then unused, since the table already carries the job's scope.
anno_relation is the small per-variant table from materialize_annotations_table, paired with a
scores-only scored_relation from materialize_scores_table. The export reconstructs scored by
LEFT-JOINing them, so each export scans the small scores-only and annotations-once tables instead of a
base that duplicates every variant's annotation blobs across all its model rows. This cuts the residual
cost below the single-table scan-once path.
export_occurrence_buckets
async
¶
export_occurrence_buckets(
*,
models: Sequence[ScoredModel],
annotation_sources: Sequence[AnnotationSource],
occurrence_relation: str,
result_relation: str,
destination_uri: str,
max_record_ordinal: int,
bucket_size: int,
rider: DisplayRider | None = None,
scored_relation: str | None = None,
anno_relation: str | None = None,
) -> list[str]
LEFT JOIN canonical results to occurrences and export one bounded Parquet prefix per ordinal bucket.
Relations and the destination are trusted operator-generated identifiers. Bucket numbers and size are validated integers. The method intentionally issues one EXPORT per bucket: BigQuery object names do not encode row ordering, while the bucket boundary gives the streaming writer a bounded sort unit.
download_header
¶
download_header(
models: Sequence[ScoredModel],
annotation_sources: Sequence[AnnotationSource],
rider: DisplayRider | None = None,
) -> list[str]
Return the ordered column names of the export_download wide TSV — the single header line to prepend
when the export is written headerless (header=False) for a server-side GCS compose. The order
matches build_export_query's SELECT aliases (a test locks this): variant_id, prioritized, the
rider's display columns, each fused annotation column, then per-model pivots named with a readable
model-name prefix, stable model-id digest, and score-column suffix.
export_prioritized
async
¶
export_prioritized(
*,
variant_ids: Sequence[str],
models: Sequence[ScoredModel],
annotation_sources: Sequence[AnnotationSource],
destination_uri: str,
rider: DisplayRider | None = None,
variants_relation: str | None = None,
scored_relation: str | None = None,
anno_relation: str | None = None,
) -> str
Export only the prioritized variants as NDJSON to destination_uri (a gs:// prefix) and return
it. This is the database branch of the two exports. Each row emits model_scores as a native
ARRAY<STRUCT>, so a loader can copy a shard into a database and decode each row with one json.loads.
That database load is the caller's step, not this method.
It differs from export_download in two ways. It keeps only prioritized variants, via a server-side
HAVING on the same prioritized rollup. It emits model_scores as a native ARRAY<STRUCT>
(ARRAY_AGG(STRUCT(…))) rather than pivoting the scores into positional columns, so the exported JSON
already has the model_scores shape the database column expects. It shares the fused scored CTE (the
join plus the pushed-down is_prioritized) with the other two export surfaces, then aggregates per
variant into the full MaterializedVariant contract:
model_scores— one struct per model that scored the variant, into_model_scorefield order (identity, score columns,COALESCE(is_prioritized, FALSE)). A coding-only variant's single model-less row is dropped (IF(model_id IS NULL, NULL, …) IGNORE NULLS), leaving an empty array.prioritized_model_ids— models whoseis_prioritizedisTRUE. The three-valuedNoneandFALSEare excluded.prioritized_source_names— each fused prioritizing source'snamewhen its flag fired, in declaration order.
Fused sources only: a cloud EXPORT runs on one backend, so unlike materialize there is no
annotate() fallback. A source that cannot sql_join raises ValueError. rider is the same optional display rider as export_download, applied to this
aggregated export and never to materialize.
This does not run on the emulator: the EXPORT DATA … OPTIONS(uri='gs://…', format='JSON') GCS sink
hangs the job poll, so the generated SQL is checked with a fake client for shape and the actual write is
an opt-in smoke test against real BigQuery. The inner SELECT (the ARRAY_AGG(STRUCT) aggregation plus
the HAVING filter) does run on the emulator and is parity-tested via build_prioritized_query.
variants_relation scopes as in materialize: a temp-table ref at warehouse scale, or None for inline
@variant_ids. scored_relation reads a pre-materialized scored table instead of rebuilding the base
CTE; pairing it with anno_relation reads the scores-only plus annotations-once split (see
export_download).
materialize_scored_table
async
¶
materialize_scored_table(
*,
dest_table: str,
variant_ids: Sequence[str],
models: Sequence[ScoredModel],
annotation_sources: Sequence[AnnotationSource],
variants_relation: str | None = None,
expiration_timestamp: str | None = None,
) -> str
Materialize the fused scored base once into dest_table with a CTAS, and return dest_table. The
join, the pushed-down is_prioritized, and the per-source flags are computed in a single
CREATE OR REPLACE TABLE … AS pass over the base scores table, which can be multi-TB. The two downloads
and the prioritized NDJSON then export off this table via scored_relation, each a small job-scoped
scan, instead of rebuilding and rescanning the base CTE. A job that emits N surfaces scans the base once
rather than N times.
The written table is exactly SELECT * FROM scored: one row per (variant, model) carrying variant_id,
model_id, model_name, the union of score_columns, every fused annotation column, is_prioritized,
and one boolean per prioritizing source. That is the full column set the scored_relation exports
project on. variants_relation and variant_ids scope the base scan as in materialize.
expiration_timestamp is a SQL TIMESTAMP expression (e.g.
TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 6 HOUR)) that sets OPTIONS(expiration_timestamp=…) so a
crashed run's temp table self-expires; None writes no expiration and the caller drops the table.
There is no ORDER BY here, unlike materialize's streaming read, so the CTAS streams rows to the table
without the global sort that pushed the prioritized export past its memory limit. This does not run on
the emulator (a CTAS over the GCS-backed base is real-BigQuery only, like EXPORT); the generated SQL is
checked with a fake client for shape, and the base CTE it wraps is the same emulator-tested
build_scored_cte.
materialize_scores_table
async
¶
materialize_scores_table(
*,
dest_table: str,
variant_ids: Sequence[str],
models: Sequence[ScoredModel],
variants_relation: str | None = None,
expiration_timestamp: str | None = None,
) -> str
Materialize the scores-only base into dest_table with a CTAS, and return dest_table. This is the
scores half of the annotations-once split. It writes one row per (variant, model) carrying only
variant_id, model_id, model_name, and the union of score_columns — no annotation columns and no
is_prioritized. The is_prioritized predicate reads region_type, a per-variant annotation, so it
can't be computed here and is deferred to the export join. This is a straight projection of the base
scores table over the job scope, with no annotation joins and no universe UNION, so the written table
holds just the numeric score rows even though the multi-TB base is scanned once.
Pair it with materialize_annotations_table, which reuses this table's variant_id set as the universe
so the base is never rescanned, and export via scored_relation=<this> and anno_relation=<that>.
variants_relation and variant_ids scope the scan as in materialize. This does not run on the
emulator (a CTAS over the GCS-backed base is real-BigQuery only); the generated SQL is checked with a
fake client for shape.
materialize_annotations_table
async
¶
materialize_annotations_table(
*,
dest_table: str,
scores_relation: str,
annotation_sources: Sequence[AnnotationSource],
variant_ids: Sequence[str],
variants_relation: str | None = None,
expiration_timestamp: str | None = None,
) -> str
Materialize the per-variant annotations-once table into dest_table with a CTAS, and return
dest_table. This is the annotations half of the split. It writes one row per universe variant carrying
every fused annotation column — the JSON blobs stored once, not duplicated across a variant's model rows
— plus one boolean per prioritizing source (see build_annotations_cte). The universe is the
already-materialized scores-only table's variant set, unioned with each prioritizing source's variants,
so this pass reads scores_relation (cheap, single-column) and each annotation source once. The base
scores table is not rescanned.
Pair it with materialize_scores_table and export via scored_relation and anno_relation.
variants_relation and variant_ids prune each prioritizing source's universe arm as in materialize.
This does not run on the emulator (a CTAS is real-BigQuery only); the generated SQL is checked with a
fake client for shape, and the two-table export join it feeds is emulator-tested via
build_export_query and build_prioritized_query with both relations set.
BigQueryDetailStore
¶
Bases: DetailStore
Generic BigQuery adapter for every binding-declared :class:DetailSchema.
table_id names the canonical JSON-row table. Its metadata sidecar is named
{table_id}__tables and exact compatible provenance is retained in {table_id}__provenance. All
schemas are fixed and contain no binding- or model-specific columns.
Writes stage one bounded batch and MERGE on (model_id, detail_name, logical_key), making
retries idempotent while preserving null-bearing logical keys through their canonical digest.
connect
classmethod
¶
connect(
table_id: str, *, credentials: Any = None
) -> BigQueryDetailStore
Build a store against a real BigQuery client. Requires the altar[bigquery] extra.
get_missing_variants
async
¶
Return requested variants without rows, reading only the clustered identity columns.
The rows table is unpartitioned and CLUSTER BY model_id, detail_name, variant_id. Each chunk
filters by equality on the first two clustering columns and IN UNNEST on the third, so BigQuery
can prune storage blocks to the requested model/table and variants instead of scanning the table.
Only those three narrow STRING columns are referenced; row_json is never read, so a wide detail
table (thousands of Enformer tracks per variant) adds rows but not payload bytes to the scan. One query
runs per 10,000 IDs; :func:~altar.scoring.prepare_scoring_batches already bounds each call by its
cache_query_size (default 10,000) and passes only primary cache hits.
DetailRowKey
dataclass
¶
Backend-neutral logical identity for one repeated result row.
digest
property
¶
Return the canonical key hash stores may use for NULL-safe uniqueness.
DetailStore
¶
Bases: ABC
Persistence axis for binding-declared named detail tables.
ensure_table
abstractmethod
async
¶
ensure_table(
*,
model_id: str,
model_name: str,
plugin_identity: PluginRunIdentity,
schema: DetailSchema,
) -> None
Create or validate a table without silently changing its schema or plugin identity.
get_missing_keys
abstractmethod
async
¶
get_missing_keys(
*,
model_id: str,
detail_name: str,
keys: Sequence[DetailRowKey],
) -> set[DetailRowKey]
Return requested logical keys not already persisted for this model and table.
get_missing_variants
async
¶
Return requested parent variants that have no persisted row in this model's table.
A detail table is one-to-many, so the logical keys a variant will produce are unknown before inference. This parent-level presence query lets cache-miss preparation treat a variant with no stored detail rows as incomplete without enumerating keys. It does not prove that a variant with some rows has all of them; callers combine it with the engine's details-before-primary write order.
The default streams :meth:read_rows and is correct for any adapter. Adapters should override it with
an indexed existence query so that wide tables (thousands of tracks per variant) are not read.
add_rows
abstractmethod
async
¶
add_rows(
*,
model_id: str,
model_name: str,
plugin_identity: PluginRunIdentity,
schema: DetailSchema,
rows: Sequence[Mapping[str, object]],
) -> InsertResult
Persist schema-shaped rows and report adapter-side duplicate skips.
read_rows
abstractmethod
async
¶
read_rows(
*,
model_id: str,
detail_name: str,
variant_ids: Sequence[str] | None = None,
) -> AsyncIterator[VariantRecord]
Stream stored rows, optionally restricted to parent variants.
describe_table
abstractmethod
async
¶
describe_table(
model_id: str, detail_name: str
) -> StoredDetailTable | None
Return durable schema/plugin provenance, or None when the table is absent.
StoredDetailTable
dataclass
¶
StoredDetailTable(
model_id: str,
model_name: str,
plugin_identity: PluginRunIdentity,
schema: DetailSchema,
)
Persisted scientific and runtime identity for one model's named detail table.
AnnotationPrioritizer
dataclass
¶
AnnotationPrioritizer(name: str, predicate: Predicate)
A reference source that can prioritize a variant without a model, such as AlphaMissense.
It pairs a name with a Predicate that decides whether the variant is prioritized, evaluated over
the variant's annotation columns. Unlike a model, it has no model_id and produces no score entry.
When its predicate is true, it only contributes to the overall prioritized flag.
annotation_prioritizers_from builds these from the configured sources.
MaterializedVariant
dataclass
¶
MaterializedVariant(
variant_id: str,
model_scores: list[dict[str, Any]],
prioritized: bool,
prioritized_model_ids: list[str],
prioritized_source_names: list[str],
plugin_identities: dict[str, dict[str, Any]],
score_lineages: dict[str, str] = dict(),
annotations: dict[str, Any] = dict(),
)
The per-variant summary produced by materialize_variant.
It holds only the fields that are the same across every model architecture. The optional display values — the headline score magnitude and the most-active cell type — are decided by each plugin and are not fields here.
annotations maps each declared annotation column that has a value for this variant to that value, merged
across every configured source. A column with no value is absent rather than None: a source with no row
for the variant, a null cell, and an empty repeated value all leave the column out (BigQuery returns a
null array as an empty one, so no store can tell those apart). A variant no source annotated gets {}.
Values keep their logical type on every store: a repeated column is a list and a struct column a dict
(a list of dicts when repeated), never the JSON text a backend such as SQLite stores them as. Values
nested inside a struct or list are kept as the source returned them.
ModelScore
dataclass
¶
ModelScore(model: ScoredModel, score: VariantRecord)
One model's raw score for a variant, keyed by the model's score_columns (for example logfc or
in_peak).
ScoredModel
dataclass
¶
ScoredModel(
model_id: str,
model_name: str,
plugin: ModelResults,
plugin_identity: PluginRunIdentity,
reuse_policy: ScoreReusePolicy = ScoreReusePolicy(),
)
A model that scored the job's variants, together with its architecture plugin.
plugin is typed as ModelResults, the smaller interface that exposes only the methods for reading
results. Materialization never needs the model-building methods (build_*), so it cannot call them by
accident. A full ModelPlugin also satisfies this type.
ScoreLineage
dataclass
¶
ScoreLineage(
lineage_id: str,
plugin_id: str,
columns: tuple[ScoreColumn, ...],
metadata_json: str = "{}",
run_identity: PluginRunIdentity | None = None,
)
An immutable catalog entry referenced by a compact ID on score rows.
An imported lineage may cover many models and input genomes: those remain
separate row keys. metadata_json holds publisher/operator declarations,
including any model-specific settings. Historical imports need no invented
Altar run identity. Registering a lineage does not by itself accept it for reuse.
from_run
classmethod
¶
from_run(
identity: PluginRunIdentity,
columns: Sequence[ScoreColumn],
) -> ScoreLineage
Key the lineage on the run's scientific projection; keep the exact identity for provenance.
The ID hashes PluginRunIdentity.scientific_dict(), so two runs share a lineage exactly when
their results are scientifically interchangeable (see same_lineage). Exact provenance (resource
locations, the declared configuration serialization, framework versions) never decides reuse and
is never repeated on score rows. A lineage store keeps one deterministic representative run identity
as the catalog's run_identity and records every distinct exact identity that registers the lineage
to write scores; ScoreStore.read_lineage_provenance returns that set.
same_lineage
¶
same_lineage(other: ScoreLineage) -> bool
Whether other declares this lineage, so a catalog holding either accepts the other.
Both must name the same ID, plugin, output columns, and metadata. A generated lineage (one with a
run_identity) is defined by its run's PluginRunIdentity.scientific_dict(), so declarations
from runs that differ only in exact provenance agree. An imported lineage has no run identity and
agrees only with an identical declaration; it never agrees with a generated one.
validate_outputs
¶
validate_outputs(
plugin_id: str, columns: Sequence[ScoreColumn]
) -> None
Acceptance cannot cross plugins or substitute missing/incompatible outputs.
ScoreReusePolicy
dataclass
¶
Additional acceptable lineages, in preference order after the current run.
This is an explicit scientific decision by the caller, scoped to the model and genome being queried. It is neither inferred from version numbers nor transitive. An empty policy accepts only the current run. Acceptance changes reads, never the lineage stamped on newly computed scores.
StoredScore
dataclass
¶
StoredScore(values: VariantRecord, lineage: ScoreLineage)
A selected score and the lineage it was stored under, including imported provenance.
In compact-lineage mode lineage.run_identity is the lineage's reference identity, one deterministic
representative of the runs that wrote under it; a row is not attributed to the exact identity that wrote
it. ScoreStore.read_lineage_provenance returns the set of exact identities that wrote under the lineage.
SchemaCollisionError
¶
Bases: ValueError
Raised when independently owned plugin columns cannot be composed safely.
ScoreStore
¶
Bases: ABC
Base class for a store that writes and reads per-model variant scores.
Stores register under altar.score_stores and are resolved by name. name is the registry key, for
example "sqlite" or "bigquery".
A note on the parameter types. At the call sites, rows (and, for a warehouse store, variant_ids) may
be pandas objects, but core itself does not depend on pandas, so they are typed as Tabular (a name for
Any documented in types); the conformance suite is what actually pins their shape. The read methods return the portable
VariantRecord instead. The columns a store persists come from plugin.score_columns(), which is
the one place the score-column set is defined.
register_score_lineage
async
¶
register_score_lineage(lineage: ScoreLineage) -> None
Register immutable metadata before importing/backfilling its compact row ID.
This does not rewrite scores or accept the lineage for any request.
Re-registering a declaration that ScoreLineage.same_lineage accepts is idempotent and keeps the
catalog entry; any other change to an existing ID fails.
read_lineage_provenance
async
¶
read_lineage_provenance(
lineage_id: str,
) -> tuple[PluginRunIdentity, ...]
Return the set of distinct exact run identities under a lineage, oldest record first.
A generated lineage is keyed on its run's scientific projection, so runs that differ only in exact
provenance (resource locations, score-neutral configuration serialization, framework API versions)
share it and reuse each other's scores. The catalog keeps one deterministic reference run_identity
per lineage, always included here; each run that registers to write scores under the lineage (even if
all its rows are skipped as duplicates) or registers its generated lineage explicitly is recorded once.
Rows are attributed to the lineage, not to one of these identities. An imported lineage, or an ID
this store never registered, has none.
get_missing_scores
async
¶
get_missing_scores(
variant_ids: Sequence[str],
model_id: str,
*,
plugin_identity: PluginRunIdentity,
columns: Sequence[ScoreColumn],
reuse_policy: ScoreReusePolicy = ScoreReusePolicy(),
) -> list[str]
Return candidate misses under a model-, genome-, and lineage-scoped policy.
read_selected_scores
async
¶
read_selected_scores(
variant_ids: Sequence[str],
model_id: str,
*,
plugin_identity: PluginRunIdentity,
columns: Sequence[ScoreColumn],
reuse_policy: ScoreReusePolicy = ScoreReusePolicy(),
) -> dict[str, StoredScore]
Select one whole row per variant, retaining its actual lineage.
Current-run results win; otherwise the first available accepted lineage wins. Output presence is a catalog declaration, not a test for non-NULL values.
required_row_identity
¶
Return typed row-identity fields this adapter requires in addition to variant_id.
The common engine preserves only explicitly declared portable identity fields. Persistence-only
stores normally need no decomposed locus and inherit the empty mapping. Warehouse stores can require
the portable chr/pos/ref/alt locus without teaching a model binding about that backend:
the engine derives omitted locus fields from the canonical variant_id and verifies emitted ones.
Portable fields keep their portable dtypes; other required fields must be emitted by the codec.
ensure_schema
abstractmethod
async
¶
ensure_schema(columns: list[ScoreColumn]) -> None
Additively set up the backing store to hold columns plus shared variant identity.
Existing model schemas remain intact. Reusing a column name is valid when its storage definition
agrees with the existing one; display-label changes do not alter stored values. Both reference stores
take the column set from score_columns().
add_scores
abstractmethod
async
¶
add_scores(
*,
model_id: str,
model_name: str,
plugin_identity: PluginRunIdentity,
rows: Tabular,
columns: list[ScoreColumn],
) -> InsertResult
Append one model's per-variant scores. rows holds the columns plus the variant-identity columns.
Returns an InsertResult: added is the number of rows this call newly persisted, and skipped
is the number the store declined because it already held a score for that (model_id, variant_id).
Idempotency is adapter-specific and NOT guaranteed by this contract. A store MAY dedup on re-add
(the sqlite store uses INSERT OR IGNORE, so re-adding a variant lands in skipped) or MAY append
unconditionally (the BigQuery store uses WRITE_APPEND, always reporting skipped=0 and
physically duplicating a re-added row). Callers therefore dedup upstream — get_unscored_variants
is the gate that keeps already-scored variants out of rows — and must not rely on add_scores
itself to be idempotent.
read_run_identity
abstractmethod
async
¶
read_run_identity(
model_id: str,
) -> PluginRunIdentity | None
Return recorded provenance compatible with model_id scores, or None when none exist.
A store may retain more than one exact run identity when compatible wrapper releases incrementally populated the model cache; this method returns one representative immutable identity. The engine compares its role-specific scientific compatibility before treating existing rows as cache hits.
get_unscored_variants
abstractmethod
async
¶
Return the variant_ids this store has no score for under model_id, so already-scored variants
are not scored again. Mirrors AnnotationSource.get_unannotated_variants, which names the same input
the same way.
read_score
async
¶
Read one model's stored score columns for one variant. Return None if this store holds no score
for that (model_id, variant_id).
This is the only read a persistence-only store has to implement. The materialize driver calls it,
through read_scores, to build each variant's list of ModelScore. The return value maps the
plugin's score_columns keys to their values, with boolean columns decoded back from whatever the
store used to save them. The method is async so a database-backed store can await its round trip
instead of blocking the event loop. The default implementation raises, so a store that overrides
materialize with a single query (such as BigQuery) does not have to provide it.
read_scores
async
¶
Read every stored score for model_id among variant_ids, keyed by variant_id. Variants this
store has no score for are left out of the returned mapping.
The materialize driver calls this rather than read_score, so it does one read per model instead
of one read per (model, variant) pair. With N variants and M models, that is M reads rather than
M×N single-row round trips on the event loop. The default implementation here just calls read_score
for each id. A store backed by a query engine should override it with a single
variant_id IN (...) fetch, as SqliteScoreStore does. A persistence-only store can implement
only read_score and inherit this.
materialize
async
¶
materialize(
*,
variant_ids: Sequence[str],
models: Sequence[ScoredModel],
annotation_sources: Sequence[AnnotationSource],
) -> AsyncIterator[MaterializedVariant]
Produce one per-variant summary for each distinct id in variant_ids, in ascending variant_id order.
This is the default driver that persistence-only stores inherit.
Every store yields the same rows in the same order: exactly one MaterializedVariant per distinct
requested variant — a repeated id is materialized once — ordered by canonical variant_id (plain
string order, as SQL ORDER BY variant_id sorts it). Request order is not kept: a warehouse store
streams its result sorted by variant_id so it can group each variant's rows without buffering the
job, and a relation-scoped request has no request order to keep. A non-canonical id raises
VariantIdentityError.
It first loads and merges each variant's annotations with fetch_annotations. Then, for each
variant, it collects that variant's per-model scores through read_scores and passes them to
materialize_variant, which does the join and prioritization the same way for every store. It yields
one MaterializedVariant per variant. A persistence-only store inherits this method unchanged and
implements only read_score.
The method is an async generator: a job with millions of variants must not hold the whole result in
memory, so callers async for over the yielded variants. Each yielded MaterializedVariant carries
model_scores, the prioritized / prioritized_model_ids / prioritized_source_names rollup, and
the variant's annotations. Every store fills annotations the same way: the declared columns that have
a value, with their logical types (see MaterializedVariant). The most-active-cell-type pick and the
headline sort value are computed by a separate plugin-declared display rollup, not here.
A warehouse store that can do the join and prioritization in one query overrides this method; see
BigQueryScoreStore. Anything specific to such a backend — for example passing a temporary-table
reference when the id list is too long for one request — lives on that override, not on this driver.
export_download
async
¶
export_download(
*,
variant_ids: Sequence[str],
models: Sequence[ScoredModel],
annotation_sources: Sequence[AnnotationSource],
destination_uri: str,
) -> str
Write all variant_ids as one flat, wide TSV to destination_uri and return that URI. Optional.
This produces a different artifact from materialize. It covers every variant, not only the
prioritized ones, and lays out each variant as one row with the per-model score columns side by side.
Its consumer is a file download (a signed URL), not the database rows. It is a separate path from
materialize and should stay that way.
The default implementation raises. This is an optimization for backends that can export in bulk, such as
BigQuery's EXPORT DATA. The sqlite store and flat-file stores do not implement it: a small
deployment can materialize and write the file itself and does not need a multi-gigabyte TSV behind a
signed URL. Anything specific to the warehouse backend, such as scoping the export to a temporary-table
reference, lives on the BigQuery override rather than in this base signature; see materialize.
SqliteScoreStore
¶
Bases: ScoreStore
Score store backed by SQLite. It uses one table, variant_scores, keyed by (model_id, variant_id),
with one column per plugin score column. New columns are added by ensure_schema without dropping any.
read_score
async
¶
Return one model's stored score for a variant, or None if there is none. Calls the bulk
read_scores so that the row-decoding logic lives in one place.
read_scores
async
¶
Return the stored scores for model_id across the given variant_ids, keyed by variant id.
This fetches them in variant_id IN (...) queries of a few hundred ids each, so the inherited
materialize driver does one bulk read per model instead of one read per (model, variant) pair, and a
long id list stays under SQLite's host-parameter limit. Bool columns are converted
back from their 0/1 storage. The model-specific schema is persisted alongside the score rows, so
ensuring another plugin's schema cannot change the shape of this read.
close
¶
Close the underlying connection. Call this for a file-backed store when you are done with it.
SqliteDetailStore
¶
Bases: DetailStore
Lossless named detail storage using schema-independent SQLite tables.
Rows are kept as canonical JSON while their parent variant and logical-key digest are indexed. This reference layout can accept any binding schema without adding physical columns or coupling a backend to a model. A provenance sidecar retains every compatible exact run identity used to populate a logical table. A production adapter may use native typed columns while preserving the same contract.
InsertResult
dataclass
¶
The result of a score or annotation write: how many rows were added and how many were skipped
because the store already held them.
Both data-lake roles return this — ScoreStore.add_scores and AnnotationSource.add_annotations — so
it lives in this shared module rather than in either role's own module. Whether a write dedups at all is
adapter-specific: an append-only store (for example the BigQuery score store) always reports
skipped=0, while a store that ignores re-adds (the sqlite stores) counts them under skipped. See
ScoreStore.add_scores for the idempotency contract.
ScoreColumn
dataclass
¶
One per-variant score column a model plugin produces.
Declaring the column here sets four things that must otherwise be kept in sync by hand: the score-store schema (table columns / DDL), the write-time validator, the wide-TSV column names, and the frontend column definitions.
It lives in this shared leaf module — not in model/plugin.py — because it is the schema the score
store consumes (ScoreStore.ensure_schema / add_scores); keeping it here lets the data-lake axis
name it without importing back into the model axis. model.plugin re-exports it for architecture
authors.
get_detail_store_registry
¶
get_detail_store_registry() -> PluginRegistry[
type[DetailStore]
]
Return the shared registry of detail-store adapters.
annotation_prioritizers_from
¶
annotation_prioritizers_from(
sources: Iterable[Any],
) -> list[AnnotationPrioritizer]
Build an AnnotationPrioritizer for each source that can prioritize a variant.
A source can prioritize when its prioritize_predicate() returns a predicate rather than None;
sources that return None only annotate and are skipped. Each source is used by duck typing — its
.name and .prioritize_predicate() — so AnnotationSource does not have to be imported here.
materialize_variant
¶
materialize_variant(
variant_id: str,
model_scores: Sequence[ModelScore],
annotation: VariantRecord,
annotation_prioritizers: Sequence[
AnnotationPrioritizer
] = (),
) -> MaterializedVariant
Combine one variant's per-model scores and annotations into a MaterializedVariant.
The steps:
- Per-model prioritization: evaluate each model plugin's predicate over the variant's scores merged
with its annotations. Keep the three-valued result (
True,False, orNone) so that a missing input yieldsNonerather than a forcedFalse. model_scores: oneto_model_scoreentry per model, with that model's three-valuedprioritizedincluded.to_model_scoreturns aNoneintoFalse.prioritized_model_ids: the ids of the models whose entry is prioritized.- Model-less prioritizers: evaluate each
AnnotationPrioritizer's predicate over the variant's annotations. A missing column makes the predicateNone, so it does not fire. The ones that do fire are collected inprioritized_source_names. prioritized: true when any model or any prioritizing source isTrue, otherwiseFalse.annotations: the columns ofannotationthat have a value (seepresent_annotations).
The most-active cell type and the headline sort magnitude are not computed here.
get_score_store_registry
¶
get_score_store_registry() -> PluginRegistry[
type[ScoreStore]
]
Return the registry of score-store plugins registered under altar.score_stores.