Skip to content

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

drop_staged_scores(staging_table: str) -> None

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, in to_model_score field 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 whose is_prioritized is TRUE. The three-valued None and FALSE are excluded.
  • prioritized_source_names — each fused prioritizing source's name when 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

BigQueryDetailStore(client: Any, table_id: str)

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

get_missing_variants(
    *,
    model_id: str,
    detail_name: str,
    variant_ids: Sequence[str],
) -> set[str]

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

DetailRowKey(
    variant_id: str, values: tuple[DetailKeyValue, ...]
)

Backend-neutral logical identity for one repeated result row.

digest property

digest: str

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

get_missing_variants(
    *,
    model_id: str,
    detail_name: str,
    variant_ids: Sequence[str],
) -> set[str]

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

ScoreReusePolicy(
    accepted_lineage_ids: tuple[str, ...] = (),
)

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

required_row_identity() -> dict[str, str]

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

get_unscored_variants(
    variant_ids: Tabular, model_id: str
) -> Tabular

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_score(
    model_id: str, variant_id: str
) -> VariantRecord | None

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_scores(
    model_id: str, variant_ids: Sequence[str]
) -> dict[str, VariantRecord]

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

SqliteScoreStore(
    path: str = ":memory:", *, lineages: bool = False
)

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

read_score(
    model_id: str, variant_id: str
) -> VariantRecord | None

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

read_scores(
    model_id: str, variant_ids: Sequence[str]
) -> dict[str, VariantRecord]

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() -> None

Close the underlying connection. Call this for a file-backed store when you are done with it.

SqliteDetailStore

SqliteDetailStore(path: str = ':memory:')

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

InsertResult(added: int, skipped: int = 0)

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

ScoreColumn(
    name: str,
    dtype: ScalarDtype,
    label: str,
    nullable: bool = True,
)

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, or None) so that a missing input yields None rather than a forced False.
  • model_scores: one to_model_score entry per model, with that model's three-valued prioritized included. to_model_score turns a None into False.
  • 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 predicate None, so it does not fire. The ones that do fire are collected in prioritized_source_names.
  • prioritized: true when any model or any prioritizing source is True, otherwise False.
  • annotations: the columns of annotation that have a value (see present_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.