Evidence sources¶
Annotation and variant–gene relation bindings should import from altar.sources.
sources
¶
Stable annotation-source and variant-to-gene evidence contracts.
This is the supported import surface for bindings backed by precomputed data, warehouses, or services. The exported records retain source evidence without forcing it into a model result.
ANNOTATION_IDENTITY_COLUMN
module-attribute
¶
The column a writable lake stamps with its source's AnnotationSourceIdentity.identity_hash.
An identity-scoped BigQueryAnnotationSource writes it on every row and deduplicates and reads only rows that
carry its own hash, so a data release bump re-annotates instead of serving cached values of the old release. A
host that loads such a lake with its own job stamps the same column. No annotation column may use this name.
VariantIdentityError
¶
Bases: ValueError
A variant key is malformed, non-canonical, or disagrees with its locus.
VariantKey
dataclass
¶
Canonical one-based biallelic key serialized as chr:pos:ref:alt.
Fields follow the rules in this module's documentation. The class normalizes spelling but does not normalize biological representation. Callers must left-align and trim indels against the selected reference before constructing a key when equivalence across representations matters.
parse
classmethod
¶
parse(variant_id: str) -> VariantKey
Parse and normalize a four-field textual key.
require_canonical
classmethod
¶
require_canonical(variant_id: str) -> VariantKey
Parse variant_id and reject aliases instead of silently rewriting them.
from_fields
classmethod
¶
from_fields(
chromosome: str,
position: int,
reference_allele: str,
alternate_allele: str,
) -> VariantKey
Construct a canonical key from decomposed one-based locus fields.
AllelicKey
dataclass
¶
AllelicKey(
variant_id: str,
chromosome: str,
position: int,
reference_allele: str,
alternate_allele: str,
)
One backend-ready biallelic variant key.
variant_id is the caller's stable identity and is returned unchanged. The remaining fields are the
values used by the backend's physical lookup; a binding may therefore add or remove a chromosome prefix
while retaining the original variant_id.
AllelicKeyColumns
dataclass
¶
Physical column names used to join a record collection to requested variants.
as_tuple
¶
Return key columns in chromosome, position, reference, alternate order.
AllelicRecord
dataclass
¶
One source-native record matched to the caller's variant_id.
AllelicRecordBackend
¶
Bases: Protocol
Retrieve source-native records without owning their biological interpretation.
fetch
async
¶
fetch(query: AllelicRecordQuery) -> Sequence[AllelicRecord]
Return all matching records, restricted to query.keys and query.fields.
Zero, one, or many records may match each variant. Returned variant_id values must come from the
request, fields must equal the requested projection, and values must not be aggregated or reinterpreted.
AllelicRecordQuery
dataclass
¶
AllelicRecordQuery(
keys: tuple[AllelicKey, ...], fields: tuple[str, ...]
)
A projected, batch lookup against one source-native record collection.
ParquetAllelicRecordBackend
¶
ParquetAllelicRecordBackend(
dataset: str | Path,
*,
key_columns: AllelicKeyColumns,
hive_partitioning: bool = True,
)
Query any Parquet record collection through DuckDB.
The adapter knows only the physical allelic key columns. Projected fields are supplied per request by the scientific binding, so one instance shape works for transcript-level AlphaMissense rows, gene-level SpliceAI rows, and future precomputed sources.
fetch
async
¶
fetch(query: AllelicRecordQuery) -> list[AllelicRecord]
Return requested native records without blocking the event loop.
SqliteAllelicRecordBackend
¶
SqliteAllelicRecordBackend(
database: str | Path,
*,
table: str,
key_columns: AllelicKeyColumns,
)
Query any SQLite table containing source-native allelic records.
A fresh connection is opened for each fetch on its worker thread, which keeps the adapter safe to call
from an async application without sharing a synchronous sqlite3.Connection across threads.
fetch
async
¶
fetch(query: AllelicRecordQuery) -> list[AllelicRecord]
Return requested native records without blocking the event loop.
AnnotationDependencyError
¶
Bases: ValueError
Raised when the configured annotation sources cannot satisfy a model's annotation dependency.
AnnotationGenomeBuildError
¶
Bases: AnnotationDependencyError
Raised when annotation sources, or sources and models, declare different genome builds.
ResolvedAnnotationDependency
dataclass
¶
ResolvedAnnotationDependency(
model_id: str,
dependency: AnnotationDependency,
source_names: tuple[str, ...],
identity: AnnotationSourceIdentity | None,
)
How one model's AnnotationDependency is satisfied.
source_names are the configured sources that supply the dependency's columns, in configuration order.
identity is the identity of the source that declares the dependency's contract, or None when the
columns were matched by name only (an UndeclaredAnnotationContractWarning was emitted).
UndeclaredAnnotationContractWarning
¶
Bases: UserWarning
Emitted when a dependency's columns are supplied by sources that do not declare its contract.
The columns are used by name, so materialization proceeds, but nothing states that their values follow the
contract the model depends on. Declare it with identity=CONTRACT.identity(release=..., genome_build=...)
on a generic source, or by overriding AnnotationSource.identity().
AnnotationColumn
dataclass
¶
AnnotationColumn(
name: str,
dtype: AnnotationDtype,
label: str,
repeated: bool = False,
fields: tuple[AnnotationColumn, ...] = (),
)
One per-variant annotation column a source provides, such as region_type or ccre_id.
dtype, together with repeated and fields, gives the column's logical type; each backend maps that
type to its own physical form. A column with repeated=True is an array of dtype. A column with
dtype="struct" is a record whose fields are declared in fields, usually with repeated=True for an
array of records such as nearest_genes. BigQuery stores these as native types (ARRAY<STRING>,
ARRAY<STRUCT<…>>). A backend without native array or record support, such as SQLite, stores them as a
JSON blob instead. That JSON form is the backend's substitute, not the column's logical type.
AnnotationContract
dataclass
¶
AnnotationContract(
source_id: str,
schema_version: str,
columns: tuple[AnnotationColumn, ...],
)
A published, versioned column set that any source may implement.
source_id is a globally namespaced name, such as org.kundajelab.altar.annotation.regions, and is what a
model's AnnotationDependency.source_id refers to. schema_version is a semantic version: a source whose
contract has the same major version and an equal or later minor version can satisfy a dependency, because
minor versions only add columns. columns is the column set every implementation declares.
A contract does not name a data release or a genome build. Several sources can implement one contract, such
as a shipped binding and a host's own warehouse table; each states the release and build it serves through
identity(release=..., genome_build=...).
identity
¶
identity(
*, release: str, genome_build: str
) -> AnnotationSourceIdentity
Return the identity of a source that serves release of this contract on genome_build.
AnnotationIdentityError
¶
Bases: ValueError
Raised when a writable annotation lake would mix, or cannot vouch for, an annotation source identity.
Examples: a row carries another identity's hash; a SQLite lake records a different identity than the source
declares; a source that is not identity-scoped writes to a scoped lake; or a lake written before identities
were stored is opened by an identity-declaring source and has not yet been adopted with
backfill_annotation_identity().
AnnotationSource
¶
Bases: ABC
A per-variant reference annotation dataset.
Adapters register under altar.annotation_sources and are looked up by name. name is the registry key,
such as "variant_annotations", "alphamissense", or "none".
annotate returns annotations in a portable shape, AnnotationTable (variant_id -> {column: value}),
which merge_annotations and fetch_annotations consume. The BigQuery source returns a pandas frame
instead and usually fuses its read into the store's query through the SqlAnnotationSource capability,
skipping the portable merge. Core carries no pandas dependency, so the signatures at this boundary use
Tabular, a documented alias for Any, and the conformance suite checks the actual behavior.
A source does not have to be registered as an entry point. A host can subclass this class, or configure a
generic BigQueryAnnotationSource or SqliteAnnotationSource over its own tables, and pass the instance
to materialization next to shipped bindings. Altar checks, joins and prioritizes on the declared
annotation_columns() and identity(), never on the class or the registry name.
annotation_columns
abstractmethod
¶
annotation_columns() -> list[AnnotationColumn]
Return the columns this source adds to each materialized variant row.
Mirrors ModelPlugin.score_columns() on the model axis: each producer names the columns it
contributes by its role (score_columns / annotation_columns).
annotate
abstractmethod
async
¶
Fetch annotations for variant_ids and return them keyed by variant_id.
The result is a Mapping[variant_id, Mapping[column, value]], the shape merge_annotations and
fetch_annotations consume. A variant this source has no row for is left out of the result; the
later left-merge fills it with {}.
identity
¶
identity() -> AnnotationSourceIdentity | None
Return the contract, data release and genome build this source serves, or None.
A source that implements a published AnnotationContract should return
CONTRACT.identity(release=..., genome_build=...), and its annotation_columns() should include the
contract's columns. Materialization resolves each model's AnnotationDependency against these
identities (see resolve_annotation_dependencies). The default None means the source declares no
contract: its columns can still satisfy a dependency by name, with an
UndeclaredAnnotationContractWarning.
add_annotations
async
¶
add_annotations(rows: Tabular) -> InsertResult
Write annotation rows into this source.
The default raises CapabilityNotSupported: most sources are read-only reference datasets that core
never writes to. Only a writable ingest source overrides this.
get_unannotated_variants
async
¶
Return the subset of variant_ids this source has no annotation row for.
A writable source uses this before ingest so it does not re-annotate variants it already has, the
same role get_unscored_variants plays for scores. The portable contract takes variant_ids in and
returns the unannotated variant_ids. A writable source that declares an identity() counts only rows
written under that identity, so a new data release reports every variant unannotated
(altar.testing.WritableAnnotationSourceContract checks this).
Warehouse-scale scoping — pruning against a temp-table relation or a chr clustering key instead of
binding millions of ids as an array parameter — is a SQL concern, not part of this backend-neutral
contract. A SQL-backed source widens this signature with those parameters on its own concrete method
(BigQueryAnnotationSource.get_unannotated_variants takes variants_relation / chromosomes), the
same way ScoreStore.get_unscored_variants keeps its variants_relation off the base and only the
BigQuery store declares it. The temp-table refs and clustering keys stay behind the SQL adapter.
The default raises CapabilityNotSupported. A precomputed reference dataset such as AlphaMissense or
CADD is never written through core, so it has nothing to deduplicate; only a writable ingest source,
paired with add_annotations, overrides this.
prioritize_predicate
¶
prioritize_predicate() -> Predicate | None
Return the predicate that gates the prioritized flag, or None.
A source that prioritizes variants, such as AlphaMissense, returns a Predicate evaluated over its
own per-variant annotation columns. A source that only contributes columns without driving
prioritization, such as region_type or a gnomAD frequency, returns None (the default).
This is the same Predicate type a ModelPlugin uses. It renders to SQL for the store's materialize
step or to Python for the reference driver, so a source can gate prioritized with no model,
model_id, or scoring job. Any aggregation from per-transcript rows down to one row per variant
happens inside annotate(); the predicate reads the already-aggregated column.
AnnotationSourceIdentity
dataclass
¶
What a concrete source's values mean: the contract it implements plus the data release and build.
source_id and schema_version name the AnnotationContract. release is the source's data release
label, such as "4.1.1-genomes"; it must change whenever the values the source returns could change.
genome_build is the assembly of the variant IDs the source answers for, spelled as in score identity
("hg38"). Surrounding whitespace is stripped from release and genome_build, and both must be
non-empty.
identity_hash is a deterministic hash of the four fields, suitable as a cache or deduplication key.
identity_hash
property
¶
Return sha256:<hex> over the canonical JSON of the four identity fields.
from_dict
classmethod
¶
from_dict(
data: Mapping[str, Any],
) -> AnnotationSourceIdentity
Rebuild an identity from to_dict() output, rejecting missing or unknown fields.
IdentityBackfill
dataclass
¶
What a writable lake's backfill changed: rows stamped, and rows deleted.
Returned by backfill_annotation_identity on BigQueryAnnotationSource and SqliteAnnotationSource, and by
BigQueryAnnotationSource.backfill_identity_columns (the locus columns).
UnscopedAnnotationReadWarning
¶
Bases: UserWarning
Emitted when a source that is not identity-scoped reads a table that has ANNOTATION_IDENTITY_COLUMN.
An identity-scoped writer stamps that column, so the table may hold rows of several identities, such as two
data releases. A source that is not scoped reads and counts all of them: a variant can be served an
arbitrary release's values, and deduplication counts a variant annotated under another release as
annotated. Give the source the identity it serves and identity_scoped=True; a reader with its own
relation= must also filter with the {identity_filter} token.
BigQueryAnnotationSource checks its table's schema for the column before it first reads or deduplicates
against the table. It does not check a read through relation=, whose tables only the SQL names. Writes
are not warned about: an unscoped write to a scoped lake raises AnnotationIdentityError.
BigQueryAnnotationSource
¶
BigQueryAnnotationSource(
client: Any,
table_id: str,
columns: Sequence[AnnotationColumn],
*,
name: str | None = None,
prioritize_predicate: Predicate | None = None,
relation: str | None = None,
genome_default: str | None = None,
identity: AnnotationSourceIdentity | None = None,
identity_scoped: bool = False,
)
Bases: SqlAnnotationSource
BigQuery-backed reference annotation source. It reads one table (project.dataset.table) keyed on
variant_id, with one column per injected AnnotationColumn. It is a prioritizing source when a
prioritize_predicate is injected; the passive/prioritizing distinction is configuration, not a subclass.
It shares the store's SQL backend, so it implements SqlAnnotationSource: the BigQuery store fuses its
sql_join JoinSpec into the single materialize query instead of round-tripping through annotate().
A source constructed with identity_scoped=True (which needs an identity) is an identity-scoped lake.
Every row it writes carries the identity's identity_hash in the _altar_annotation_identity column
(ANNOTATION_IDENTITY_COLUMN). Deduplication counts only rows of that hash, and reads select only rows of
that hash: the bare table through a filtering subquery, and a relation= through its required
{identity_filter} token. Bumping the identity's release therefore re-annotates variants instead of
serving the old release's cached values, and a table holding several releases still yields one row per
variant. Rows written before the table carried an identity have a NULL hash and are neither served nor counted
until backfill_annotation_identity() adopts them. Without identity_scoped, identity is metadata only.
Every source object that writes or reads the lake must be scoped the same way. A source that is not scoped
looks up its table's schema, a metadata request that bills nothing, and treats a table with the identity
column as a scoped lake. add_annotations and add_annotations_frame then raise AnnotationIdentityError
instead of writing rows without a hash. The source emits UnscopedAnnotationReadWarning, once per source
object, before it first reads the table (annotate and the store's fused queries) or deduplicates against it
(get_unannotated_variants and export_unannotated). A read through relation= is not checked, because only
its SQL names the tables it reads; deduplication and writes always use the bare table, so they are checked. If
the lookup fails, the failure is logged and the operation proceeds unchecked.
client is a google.cloud.bigquery.Client (or a compatible fake). table_id is
"project.dataset.table". columns and prioritize_predicate are the wiring that makes this a passive
or prioritizing source. name is the logical source identity recorded in prioritized_source_names
(e.g. "alphamissense"), which is distinct from the class registry key ("bigquery") and defaults to
it. The client is injected so this module needs no google-cloud-bigquery at import.
relation is an optional SQL relation that overrides the bare table_id: a subquery yielding one row
per variant_id and exposing at least the declared columns. It is how a source whose backing rows are
per-transcript (AlphaMissense and REVEL collapse to the most-damaging transcript) or that needs computed
or renamed columns (gnomAD af → af_global, TO_JSON_STRING(...) blobs) fuses in-SQL, keeping the
aggregation and projection inside the source rather than the store. A {variant_filter} token in the
string expands, in both sql_join and annotate, to variant_id IN UNNEST(@<param>), so the reference
table is pruned to the job's variants before aggregating. A {identity_filter} token expands to
_altar_annotation_identity = '<identity_hash>', so a relation over a lake holding several releases
reads only this source's rows. It qualifies like {variant_filter} does (a.{identity_filter}). An
identity-scoped source's relation must contain it, and any other source's relation must not. None uses
the bare per-variant table.
genome_default is the genome build of a writable ingest lake: the genome that add_annotations and
add_annotations_frame stamp on rows that carry none (a row naming another build is rejected), and that
both get_unannotated_variants paths and export_unannotated filter on, so the dedup only counts rows
of this build and cluster-prunes the annotation table's [genome, chr, variant_id] key. No build is
assumed: those methods raise ValueError when it is None, rather than silently deduplicating an hg19
lake against hg38 rows. It is inert for a read-only reference source, which is never written or
deduplicated.
identity declares the AnnotationContract this table implements and the data release and genome build
it holds (see AnnotationSource.identity()); a host table that implements a published contract passes
CONTRACT.identity(release=..., genome_build=...). On its own it is metadata: it does not change how rows
are written, deduplicated or read. identity.genome_build is the build label that dependency resolution
compares with score runs (for example "hg38"), while genome_default is how this table spells the
build in its stored genome column (Altar's VCF ingest writes "GRCh38"). The two are independent, and
neither is derived from the other.
identity_scoped=True makes the identity decide which rows this source writes, counts and reads (see the
class docstring). It needs identity, and a relation must then contain {identity_filter}; both are
checked here and raise ValueError. Writes stamp identity.identity_hash, and a row that carries
another hash raises AnnotationIdentityError.
connect
classmethod
¶
connect(
table_id: str,
columns: Sequence[AnnotationColumn],
*,
name: str | None = None,
prioritize_predicate: Predicate | None = None,
relation: str | None = None,
genome_default: str | None = None,
credentials: Any = None,
identity: AnnotationSourceIdentity | None = None,
identity_scoped: bool = False,
) -> BigQueryAnnotationSource
Build a source against a real BigQuery client. Requires the altar[bigquery] extra.
sql_join
¶
sql_join(ctx: SqlJoinContext) -> JoinSpec
Fuse this source into the store's materialize query as a LEFT JOIN on ctx.join_key. A BigQuery
source always shares the store's backend, so it always fuses and never returns None. relation is the
backtick-quoted table id for a per-variant table (the store projects each injected column), or the
injected aggregating or computed subquery (see __init__'s relation and JoinSpec), pruned to
ctx.variants_relation (temp table) or ctx.variant_ids_param (inline param). The alias derives from
the logical source name (alphamissense, cadd, …), which the store keeps distinct across sources.
ensure_schema
async
¶
Reconcile a writable ingest table to hold the locus identity plus the injected columns.
The identity is the variant_id key, genome, and the decomposed chr/pos/ref/alt locus that
add_annotations and add_annotations_frame both write and the dedup reads filter on, clustered by
[genome, chr, variant_id]. CREATE covers a fresh table; ADD COLUMN IF NOT EXISTS reconciles a
pre-existing one by adding any missing columns (e.g. a table provisioned out of band). BigQuery cannot
add a NOT NULL column, so identity columns added to an old table are nullable. Its earlier rows have
no genome, so the dedup reports them unannotated, and because writes append, re-annotating them would
leave a second row per variant. Run backfill_identity_columns once after this migration, before any
dedup; ensure_schema never rewrites existing rows itself.
An identity-scoped source (see the class docstring) also adds a nullable _altar_annotation_identity
STRING column. Added to an existing table, it is NULL on every earlier row, so those rows are neither
served nor counted as annotated until backfill_annotation_identity adopts them. Without that
backfill, the next deduplication re-annotates the whole lake. The clustering is unchanged.
backfill_identity_columns
async
¶
backfill_identity_columns() -> IdentityBackfill
Stamp the locus identity on rows written before the table carried it, as one transaction.
This is the explicit, one-time migration for a table that predates the identity columns (see
ensure_schema, which must run first to add them; it never runs this). It assumes every legacy row, one
with a NULL genome, belongs to genome_default's build, which is therefore required. Only a legacy row
whose variant_id is canonical (_CANONICAL_VARIANT_ID_PATTERN, the RE2 form of
VariantKey.require_canonical) is touched:
- If its variant already has a row of
genome_default's build (it was re-annotated before this ran), the legacy row is deleted. - Otherwise exactly one legacy row per
variant_idis kept, stamped withgenome_defaultand thechr/pos/ref/altparsed from thevariant_id, and the others are deleted. The survivor is the most recent bycreated_atwhen the table has that column, then the first by the row's JSON text, so a rerun on the same rows keeps the same one.
A legacy row with a non-canonical variant_id (1:9:A:T, chr6:-3:A:T) keeps its NULL genome: no
canonical lookup can reach it, and stamping it would invent a locus. The work runs as one BigQuery
multi-statement transaction (temp-table snapshot, DELETE, INSERT), so a concurrent write cannot
interleave with it; stop ingest into the table while it runs and resume afterwards. A second run finds no
canonical legacy rows and changes nothing. It returns how many rows were stamped and deleted.
backfill_annotation_identity
async
¶
backfill_annotation_identity() -> IdentityBackfill
Adopt the rows written before this lake stored an annotation identity, as one transaction.
This is an explicit, reviewed one-time migration, not something Altar runs on its own. It asserts that
every legacy row of genome_default's build, one whose _altar_annotation_identity is NULL, was produced
by this source's configured identity (the same contract, data release and genome build). Run it only
with the identity of the release that actually built those rows. Adopting them under a newer release
would serve old values under the new identity, which is the error identity scoping exists to prevent.
The source must be identity-scoped (identity_scoped=True) and have a genome_default, and
ensure_schema must have added the column. Rows with a NULL genome are not touched: run
backfill_identity_columns first for those. Of the legacy rows of this build:
- If its variant already has a row stamped with this identity (it was re-annotated before this ran), the legacy row is deleted.
- Otherwise exactly one legacy row per
variant_idis kept and stamped with the identity hash, and the others are deleted. The survivor is chosen as inbackfill_identity_columns: the most recent bycreated_atwhen the table has that column, then the first by the row's JSON text.
Legacy rows of another genome and rows stamped with any identity are left alone. Like
backfill_identity_columns it runs as one BigQuery multi-statement transaction (temp-table snapshot,
DELETE, INSERT); stop ingest into the table while it runs. A second run finds no legacy rows of this
build and changes nothing. It returns how many rows were stamped and deleted.
Cost: the snapshot reads every column of the legacy rows, and on-demand pricing bills the DELETE for
every column of the whole table, so when most rows are legacy the transaction bills about three times the
table's logical size and rewrites all of its storage. BigQuery re-clusters the rewritten rows in the
background at no charge. The replaced storage stays in time travel and fail-safe for their configured
windows, which a dataset on physical storage billing pays for.
add_annotations
async
¶
add_annotations(rows: Tabular) -> InsertResult
Append per-variant annotation rows, each a {"variant_id": ..., <column>: value, ...} mapping.
The write is append-only (WRITE_APPEND), so deduping is the caller's job (pre-flight via
get_unannotated_variants). Like add_annotations_frame, every row is written with its genome and
decomposed locus, so both write paths fill the identity columns the dedup reads filter on. variant_id
must be canonical: a non-canonical one raises VariantIdentityError rather than being written where no
read can find it. chr/pos/ref/alt are derived from it, and any the row carries must name the same
variant. genome is the row's own, else genome_default (see __init__). An identity-scoped lake also
stamps _altar_annotation_identity with its identity hash; a row that carries another hash raises
AnnotationIdentityError. The table must already have that column (ensure_schema). A source that is not
identity-scoped raises AnnotationIdentityError when its table has that column (see the class docstring).
add_annotations_frame
¶
add_annotations_frame(frame: Any) -> InsertResult
Ingest a per-variant annotation frame columnar, the annotation counterpart of add_scores_frame.
It re-derives the variant_id (chr:pos:ref:alt) join key from the locus, overwriting any
present-but-null passthrough id and rejecting a non-canonical locus with VariantIdentityError, stamps
genome_default on a frame without a genome column (and rejects a genome column naming another build
when genome_default is set), stamps an identity-scoped lake's identity hash (rejecting a frame whose
_altar_annotation_identity names another), stamps created_at, and casts each declared scalar
column to its physical dtype (nullable Int64, bool, or float64). repeated, struct, str, and
json columns pass through untouched: BigQuery's load_table_from_dataframe (Parquet) maps a
list-of-scalar to ARRAY<…> and a list-of-dict to ARRAY<STRUCT<…>>. It drops any af_* columns the
caller's frame carries, since allele-frequency columns live in a separate source joined at materialize,
so ALLOW_FIELD_ADDITION can't resurrect them. The write is append-only (WRITE_APPEND), so deduping is
the caller's job (pre-flight via get_unannotated_variants), as in the score store. A source that is not
identity-scoped raises AnnotationIdentityError when its table has _altar_annotation_identity (see the
class docstring).
get_unannotated_variants
async
¶
get_unannotated_variants(
variant_ids: Sequence[str],
*,
variants_relation: str | None = None,
chromosomes: Sequence[str] | None = None,
) -> list[str]
Return the subset of variant_ids this lake has no annotation row for. This is the ingest dedup that
lets an annotation job annotate only cache misses, the annotation counterpart of get_unscored_variants,
differing only in the missing model_id (annotations are not per-model). It binds @variant_ids inline
by default. 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 unannotated rows are read straight off
a LEFT JOIN … IS NULL dedup and variant_ids may be empty. chromosomes then prunes the annotation
table's chr cluster key; it is ignored on the inline path, whose small array binds directly. Both paths
count only rows of genome_default's build, so it is required, and an identity-scoped lake counts only
rows stamped with its identity hash, so a release bump reports every variant unannotated. It returns
variant_ids, and the caller decomposes them back to loci (chr:pos:ref:alt) for the annotation-job
input file.
export_unannotated
async
¶
export_unannotated(
*,
destination_uri: str,
variants_relation: str,
chromosomes: Sequence[str] | None = None,
) -> int
Export unannotated loci from a job relation directly to a headerless object-store TSV.
It counts and exports the same cache misses as get_unannotated_variants with variants_relation: only
rows of genome_default's build, and, for an identity-scoped lake, of its identity hash, are annotated.
annotate
async
¶
Return {variant_id: {column: value}} for the requested variants this table has rows for. A variant
with no row is absent, and the left-merge then yields {} for it (see merge_annotations). A declared
scalar json column that the relation returns as JSON text, such as TO_JSON_STRING(...), is decoded to
its JSON value, as SqliteAnnotationSource decodes its stored text. An identity-scoped lake reads only
rows of its identity, through the same relation sql_join uses.
BigQueryAllelicRecordBackend
¶
BigQueryAllelicRecordBackend(
client: Any,
table_id: str,
*,
key_columns: AllelicKeyColumns,
filters: Mapping[str, str | int | bool] | None = None,
batch_size: int = DEFAULT_FETCH_BATCH_SIZE,
)
Read source-native allelic records from one BigQuery table.
table_id is a project.dataset.table reference. key_columns names the table's chromosome, position,
reference-allele, and alternate-allele columns. filters limits every read to rows whose columns equal
fixed values. Use it when one table holds several genome builds or releases, for example
{"genome": "hg38"}. Filter values are bound as query parameters.
Each requested key must describe the locus its variant_id names. A binding may spell the chromosome
differently from the variant ID, for example 1 for chr1, but the position and alleles must match.
fetch raises ValueError for a key that does not. covered_variant_ids relies on the same rule to find
candidate variants in SQL without calling the binding.
A fetch query reads the key and requested columns of the whole table, because BigQuery can skip storage
only by clustering or partitioning. fetch therefore sends keys in batches of at most batch_size, one query
per batch. Variants returned by the most recent covered_variant_ids call are answered from memory instead.
fetch
async
¶
fetch(query: AllelicRecordQuery) -> list[AllelicRecord]
Return every record matching query.keys, with only query.fields, in request order.
covered_variant_ids
async
¶
Return the variants in variants_relation that have at least one record in this table.
variants_relation is a SQL relation with a variant_id column in chr:pos:ref:alt form, such as a
backtick-quoted table id. The match compares positions exactly and alleles ignoring case and surrounding
spaces. It compares chromosomes after removing a chr prefix and leading zeros and treating MT as M.
That comparison is looser than fetch, so it can include a variant fetch finds nothing for. It never
leaves out a variant fetch would find a record for.
The query also reads every column of the matched records and keeps them in memory. The next fetch for
a covered variant is answered from those records, with the same exact key match a query would apply, and
then releases them. If a covered variant's records are spelled differently from the requested key, for
example 1 in the table and chr1 in the key, that fetch raises ValueError rather than return no
records. A later call to this method replaces records that were never fetched. Point the backend at a
table or view without unused wide columns to keep this read small.
JoinSpec
dataclass
¶
One LEFT JOIN the store adds to its materialize query to read an annotation source in SQL.
Adding this join lets the store read the source within its single query instead of calling annotate()
and merging in Python. For each spec the store builds:
`LEFT JOIN <relation> AS <alias> ON <alias>.<join_key> = <base_alias>.<join_key>`
and adds <alias>.<column> for each entry in columns to the materialized row.
relation is whatever goes between LEFT JOIN and AS. It is a backtick-quoted table reference for a
per-variant table, or a parenthesized subquery such as (SELECT variant_id, … FROM t GROUP BY
variant_id) when the source first collapses several per-transcript rows into one row per variant
(AlphaMissense and REVEL do this inside the source, not the store). A source that cannot fuse returns
None from sql_join, and the store reads it through annotate() instead.
SqlAnnotationSource
¶
Bases: AnnotationSource
An AnnotationSource that shares a SQL backend with the store and can fuse its read into the store's
materialize query.
A source subclasses this to signal that the store may read it inside one SQL query rather than through
annotate() and a Python merge. A SQL store finds these sources with
isinstance(source, SqlAnnotationSource). A plain AnnotationSource (Parquet, in-process, or a
different backend) does not subclass it, so the store reads that source through annotate(). Subclassing
is how a source opts in; the base AnnotationSource contract has no SQL methods.
sql_join
abstractmethod
¶
sql_join(ctx: SqlJoinContext) -> JoinSpec | None
Return a JoinSpec the store adds to its materialize query, or None if this source cannot fuse
for the given ctx. When it returns None, the store reads this source through annotate() instead.
SqlJoinContext
dataclass
¶
SqlJoinContext(
base_alias: str,
join_key: str = "variant_id",
variant_ids_param: str = "variant_ids",
variants_relation: str | None = None,
)
The information the store passes each source at materialize time so it can build a join (sql_join).
It holds only what the store's materialize step currently needs.
base_alias— the store's base relation, the scores side, which each source's join attaches to as the left side.join_key— the column both sides join on (variant_id).variant_ids_param— the name of the array query parameter the store binds to the job's variant ids (@variant_ids). A source whoserelationis an aggregating subquery, where AlphaMissense or REVEL collapse several per-transcript rows into one, uses this to restrict the large reference table to the job's variants before aggregating, with anINNER JOIN variantsscan. The outerLEFT JOINcannot do this, because it runs after theGROUP BYand so cannot prune the group's input.variants_relation— an optional SQL relation (a table id or subquery) that yields the job'svariant_ids. When set,{variant_filter}prunes withvariant_id IN (SELECT variant_id FROM <relation>)instead of the inlineUNNEST(@variant_ids)parameter. A job with millions of variants uses this to scope the query without exceeding BigQuery's request-size limit: it uploads a temp table and passes the reference here. The default,None, uses the inline-parameter prune, which suits small jobs and the emulator.
SqliteAnnotationSource
¶
SqliteAnnotationSource(
path: str = ":memory:",
columns: Sequence[AnnotationColumn] = (),
*,
name: str | None = None,
prioritize_predicate: Predicate | None = None,
identity: AnnotationSourceIdentity | None = None,
identity_scoped: bool = False,
)
Bases: AnnotationSource
Annotation source backed by SQLite. It uses one table, variant_annotations, keyed by variant_id,
with one column per AnnotationColumn. It acts as a prioritizing source when you pass a
prioritize_predicate, and a passive one otherwise. That choice is configuration, not a separate
subclass.
Unlike the read-only base class, this source can also be written to. Use add_annotations to load your
own reference data so it feeds into materialize, and get_unannotated_variants to find the variants it
has no row for yet.
identity= alone is metadata. With identity_scoped=True, the file holds the annotations of exactly one
AnnotationSourceIdentity, recorded in its _altar_annotation_identity table. A scoped source records its
identity the first time it opens a file with no annotation rows. Opening a file that records a different
identity raises AnnotationIdentityError, so a data release bump uses a new file rather than serving the
old release's rows. A file that already holds rows but records no identity was written before identities
were stored: a scoped source opens it, but refuses to read or write it until backfill_annotation_identity()
adopts its rows. Each operation re-reads the recorded identity, so an adoption by another instance is seen.
An unscoped source reads any file, but cannot write to one that records an identity.
backfill_annotation_identity
async
¶
backfill_annotation_identity() -> IdentityBackfill
Adopt a legacy file's rows as this source's identity, and return how many rows were adopted.
This is an explicit, reviewed one-time migration. It records the configured identity in a file that
holds annotation rows but no identity, asserting that the declared data release produced every row in
it. It needs identity_scoped=True. It returns IdentityBackfill(stamped=<rows adopted>, deleted=0):
the primary key already allows one row per variant, so nothing is deleted. A file that already records
this identity is left unchanged and stamped is 0, so a second call is a no-op; a file recording
another identity cannot be opened at all.
add_annotations
async
¶
add_annotations(rows: Tabular) -> InsertResult
Insert per-variant annotation rows. Each row is a mapping like
{"variant_id": ..., <column>: value, ...}. Rows use INSERT OR IGNORE on the variant_id primary
key, so a variant that already has a row is skipped rather than overwritten.
For an identity-scoped source, the rows belong to the file's identity, and a row whose
_altar_annotation_identity names another identity hash raises AnnotationIdentityError. An unscoped
source cannot write to a file that records an identity.
annotate
async
¶
Return {variant_id: {column: value}} for the requested variants that this source has rows for. A
variant with no row is left out of the result; the merge step then treats it as an empty annotation
(see merge_annotations).
get_unannotated_variants
async
¶
Return the variant_ids this file has no annotation row for, in input order.
For an identity-scoped source the file holds one identity (see the class docstring), so every stored row is a row of this source's identity.
close
¶
Close the underlying connection. Call this for a file-backed source when you are done with it.
BiosampleContext
dataclass
¶
The cell, tissue, or other biosample in which a link was measured or predicted.
GenomeBuild
¶
Bases: StrEnum
Genome assemblies supported by the normalized link contract.
GenomeBuildMismatchError
¶
Bases: VariantGeneLinkQueryError
The query and source use different genome assemblies.
GenomicInterval
dataclass
¶
GenomicInterval(
chromosome: str,
start: int,
end: int,
genome_build: GenomeBuild,
coordinate_system: CoordinateSystem = "0-based-half-open",
)
A normalized 0-based, half-open genomic interval [start, end).
overlaps
¶
overlaps(other: GenomicInterval) -> bool
Return whether two same-build intervals overlap.
InvalidCursorError
¶
Bases: VariantGeneLinkQueryError
A pagination cursor is malformed or belongs to a different query.
LinkDistances
dataclass
¶
LinkDistances(
element_to_gene_bp: float | None = None,
variant_to_gene_bp: float | None = None,
variant_to_element_bp: float | None = None,
interpretation: str | None = None,
)
Distances retained when a source provides them, in base pairs.
LinkProvenance
dataclass
¶
LinkProvenance(
source_record_id: str,
source_dataset: str,
source_url: str | None = None,
study_id: str | None = None,
publication_ids: tuple[str, ...] = (),
metadata: Mapping[str, Any] = dict(),
)
Evidence identifiers and source-specific metadata retained without flattening.
LinkScore
dataclass
¶
LinkScore(
value: float | None,
direction: ScoreDirection,
interpretation: str,
threshold: float | None = None,
threshold_rule: ThresholdRule | None = None,
threshold_applied: bool | None = None,
threshold_interpretation: str | None = None,
)
A source score without conflating a missing value with numeric zero.
RegulatoryElement
dataclass
¶
RegulatoryElement(
id: str,
interval: GenomicInterval,
element_type: str | None = None,
)
Identity and genomic extent of one regulatory element.
ScoreDirection
¶
Bases: StrEnum
How to interpret increasing values of a source's link score.
TargetGene
dataclass
¶
TargetGene(
id: str,
namespace: str,
namespace_version: str | None = None,
symbol: str | None = None,
)
A target gene identifier with an explicit namespace and optional release/symbol.
VariantGeneLink
dataclass
¶
VariantGeneLink(
record_id: str,
locus: VariantLocus,
regulatory_element: RegulatoryElement | None,
target_gene: TargetGene,
context: BiosampleContext | None,
source_id: str,
source_name: str,
method: str,
data_release: str,
score: LinkScore,
distances: LinkDistances,
provenance: LinkProvenance,
)
One loss-minimizing variant/locus-to-gene evidence record.
regulatory_element is present for mediated evidence such as enhancer-to-gene links. It is absent for
sources that directly associate a locus or allele with a gene. Consumers that specifically project an
enhancer graph must reject direct records instead of inventing a regulatory element.
VariantGeneLinkDataError
¶
Bases: VariantGeneLinkSourceError
The source returned malformed, unstable, or ambiguously duplicated data.
VariantGeneLinkError
¶
Bases: Exception
Base exception for variant-to-gene link operations.
VariantGeneLinkPage
dataclass
¶
VariantGeneLinkPage(
links: tuple[VariantGeneLink, ...],
next_cursor: str | None = None,
)
A bounded page of links; an empty tuple with no cursor is a successful empty result.
VariantGeneLinkQuery
dataclass
¶
VariantGeneLinkQuery(
locus: VariantLocus,
context_ids: tuple[str, ...] = (),
context_names: tuple[str, ...] = (),
source_ids: tuple[str, ...] = (),
gene_ids: tuple[str, ...] = (),
methods: tuple[str, ...] = (),
data_releases: tuple[str, ...] = (),
minimum_score: float | None = None,
maximum_score: float | None = None,
page_size: int = 100,
cursor: str | None = None,
)
One interval-overlap query plus relation-level filters and a bounded page size.
Context IDs and names are alternative selectors (their union). Other non-empty identifier tuples are intersections, as are the inclusive score bounds. No score cutoff is applied when both bounds are absent.
VariantGeneLinkQueryError
¶
Bases: VariantGeneLinkError, ValueError
The normalized query is invalid for the selected source.
VariantGeneLinkSource
¶
Bases: ABC
Queryable one-to-many variant/locus-to-gene evidence.
A valid query with no matching links returns VariantGeneLinkPage(links=(), next_cursor=None).
Network, storage, malformed-source-data, and pagination failures raise VariantGeneLinkSourceError
(or a subclass) and must never be translated into an empty page.
query_links
abstractmethod
async
¶
query_links(
query: VariantGeneLinkQuery,
) -> VariantGeneLinkPage
Return one stable, bounded page of links for query.
stream_links
async
¶
stream_links(
query: VariantGeneLinkQuery,
) -> AsyncIterator[VariantGeneLink]
Stream all pages without retaining the full relation in memory.
VariantGeneLinkSourceError
¶
Bases: VariantGeneLinkError
A source could not answer a valid query because of network, storage, or source-data failure.
VariantLocus
dataclass
¶
VariantLocus(
interval: GenomicInterval,
reference_allele: str | None = None,
alternate_allele: str | None = None,
)
A normalized query locus, optionally carrying normalized alleles.
interval always uses 0-based, half-open coordinates. from_vcf is the explicit conversion from
a 1-based VCF/Open Targets variant position; the reference allele determines the reference span.
from_vcf
classmethod
¶
from_vcf(
*,
chromosome: str,
position: int,
reference_allele: str,
alternate_allele: str,
genome_build: GenomeBuild,
) -> VariantLocus
Convert a 1-based VCF locus into the contract's 0-based, half-open interval.
InMemoryVariantGeneLinkStore
¶
Deterministic reference store used by examples and Cartesian conformance tests.
LinkRecordWriteResult
dataclass
¶
Counts returned by an idempotent generation write.
SqliteVariantGeneLinkStore
¶
Portable SQLite store for canonical evidence with indexed overlap and relation filters.
StoredVariantGeneLinkSource
¶
StoredVariantGeneLinkSource(
store: VariantGeneLinkStore,
dataset: VariantGeneLinkDatasetRef,
*,
name: str,
binding_id: str,
binding_version: str,
failure_label: str = "variant-gene link store",
)
Bases: VariantGeneLinkSource
Reusable public source over a generic canonical evidence store.
decode_record
¶
decode_record(
record: VariantGeneEvidenceRecord, locus: VariantLocus
) -> VariantGeneLink
Project a canonical record into a result; bindings may restore typed source metadata.
VariantGeneEvidenceRecord
dataclass
¶
VariantGeneEvidenceRecord(
record_id: str,
evidence_interval: GenomicInterval,
regulatory_element: RegulatoryElement | None,
target_gene: TargetGene,
context: BiosampleContext | None,
source_id: str,
source_name: str,
method: str,
data_release: str,
score: LinkScore,
distances: LinkDistances,
provenance: LinkProvenance,
metadata: Mapping[str, JsonValue] = dict(),
)
One durable, query-independent locus/element-to-gene evidence record.
to_link
¶
to_link(locus: VariantLocus) -> VariantGeneLink
Attach a caller's overlapping locus and return the public query result.
VariantGeneLinkDatasetRef
dataclass
¶
Stable identity for one independently versioned collection of canonical link evidence.
key
property
¶
Return a deterministic store key without assigning semantics to separators in components.
VariantGeneLinkStore
¶
Bases: Protocol
Persist canonical evidence independently of its scientific source format.
Records are written into immutable generations and become queryable only after publish_generation.
This lets importers checkpoint bounded batches without exposing partial releases.
write_records
async
¶
write_records(
dataset: VariantGeneLinkDatasetRef,
generation_id: str,
records: Sequence[VariantGeneEvidenceRecord],
) -> LinkRecordWriteResult
Idempotently stage a bounded batch in generation_id.
publish_generation
async
¶
publish_generation(
dataset: VariantGeneLinkDatasetRef, generation_id: str
) -> None
Atomically make all records in a staged generation visible to readers.
generation_size
async
¶
generation_size(
dataset: VariantGeneLinkDatasetRef, generation_id: str
) -> int
Return the number of unique staged or published records in a generation.
available_builds
async
¶
available_builds(
dataset: VariantGeneLinkDatasetRef,
) -> Sequence[GenomeBuild]
Return genome builds represented by published generations.
fetch_page
async
¶
fetch_page(
dataset: VariantGeneLinkDatasetRef,
query: VariantGeneLinkStoreQuery,
) -> Sequence[VariantGeneEvidenceRecord]
Return published records ordered strictly by record_id.
VariantGeneLinkStoreQuery
dataclass
¶
VariantGeneLinkStoreQuery(
interval: GenomicInterval,
context_ids: tuple[str, ...] = (),
context_names: tuple[str, ...] = (),
source_ids: tuple[str, ...] = (),
gene_ids: tuple[str, ...] = (),
methods: tuple[str, ...] = (),
data_releases: tuple[str, ...] = (),
minimum_score: float | None = None,
maximum_score: float | None = None,
after_record_id: str | None = None,
limit: int = 101,
)
Backend-neutral indexed request derived from a public link query.
from_query
classmethod
¶
from_query(
query: VariantGeneLinkQuery,
) -> VariantGeneLinkStoreQuery
Translate the public query while decoding its opaque cursor once in the source layer.
canonical_chromosome
¶
Return the canonical chr-prefixed spelling used in portable keys.
Primary chromosomes accept common aliases (1/chr1/CHR01 and
M/MT/chrM) in any case. Other reference contigs keep their exact
spelling after an optional case-insensitive chr prefix is normalized, so
CHRUn_KI270302v1 becomes chrUn_KI270302v1 but chrun_KI270302v1
stays as written. Non-ASCII text, whitespace inside the name, and characters
outside the VCF contig-name set (including :) raise
VariantIdentityError. This is a naming rule only; it does not claim that
contigs from different assemblies are interchangeable.
canonical_variant_id
¶
canonical_variant_id(
chromosome: str,
position: int,
reference_allele: str,
alternate_allele: str,
) -> str
Return the canonical text key for decomposed one-based locus fields.
resolve_annotation_dependencies
¶
resolve_annotation_dependencies(
models: Sequence[ScoredModel],
sources: Sequence[AnnotationSource],
*,
contracts: Sequence[AnnotationContract] = (),
) -> tuple[ResolvedAnnotationDependency, ...]
Check that sources satisfy every annotation dependency of models, and return how.
For each AnnotationDependency in each model's manifest.prioritization:
- If a source's
identity().source_idis the dependency'ssource_id, that source must implement a compatible contract version (same major version, and a version no older than the dependency's) and declare every dependency column. Each of those columns must have the logical type (dtype,repeatedand nestedfields; labels may differ) that every known contract with the dependency'ssource_idand major version gives it. OtherwiseAnnotationDependencyErroris raised. - If no source declares that contract but every dependency column is declared by some source, the
dependency resolves by column name and an
UndeclaredAnnotationContractWarningis emitted. - Otherwise
AnnotationDependencyErroris raised naming the contract and the missing columns.
Every source that declares an identity must declare the same genome_build, and a model whose run
identity records a genome_build must match it; otherwise AnnotationGenomeBuildError is raised.
The known contracts are the ones Altar publishes, such as REGION_ANNOTATION_CONTRACT, plus contracts. A
dependency on a contract the resolver does not know, such as a binding's or a host's own, gets the name check
only. Pass that contract in contracts to check its column types too, or run
altar.testing.assert_implements_contract to check a source against its whole contract.
Store materialize methods, and the BigQuery exports that compute prioritized, call this once per call,
without contracts. A host can also call it up front, with its own contracts, to check a configuration
before submitting work.
fetch_annotations
async
¶
fetch_annotations(
sources: Sequence[AnnotationSource],
variant_ids: Sequence[str],
) -> dict[str, dict[str, Any]]
Fetch each source's annotations for variant_ids and merge them with merge_annotations.
This builds the annotation input to materialize_variant from the configured annotation sources.
Sources are fetched concurrently and merged in declaration order. Their declared column names must be
globally unique; a collision is rejected before any source is called. Each annotate must return the
portable variant_id -> {column: value} mapping.
get_annotation_source_registry
¶
get_annotation_source_registry() -> PluginRegistry[
type[AnnotationSource]
]
Return the shared registry of annotation sources (the altar.annotation_sources group).
merge_annotations
¶
merge_annotations(
variant_ids: Sequence[str],
tables: Sequence[AnnotationTable],
*,
source_names: Sequence[str] | None = None,
) -> dict[str, dict[str, Any]]
Left-merge several sources' annotation tables into one mapping per variant.
The left side is variant_ids, the variants being scored. Every id gets an entry, even one that no source
annotated, which yields {}. Each table is the portable variant_id -> {column: value} shape a source
returns from annotate. Two sources may not supply the same column: doing so is ambiguous even when
their declared dtypes agree, because their values may differ. A row a source returns for a variant not
in variant_ids is dropped.
This is the portable join that builds the annotation argument to materialize_variant. The BigQuery
store instead fuses its annotation reads into one query through SqlAnnotationSource.sql_join.
staged_annotation_source
async
¶
staged_annotation_source(
source: AnnotationSource,
client: Any,
*,
staging_dataset: str,
variants_relation: str | None = None,
variant_ids: Sequence[str] | None = None,
coverage: BigQueryAllelicRecordBackend | None = None,
batch_size: int = DEFAULT_STAGE_BATCH_SIZE,
expiration: timedelta | None = DEFAULT_STAGE_EXPIRATION,
) -> AsyncIterator[BigQueryAnnotationSource]
Stage source's annotations for one job in a scratch table and yield a SQL source bound to it.
source is a source without SQL of its own. The yielded source has the same name, columns, identity, and
prioritization rule, and reads the scratch table. Its joins are limited to each store query's variants.
Give the job's variants as exactly one of variant_ids or variants_relation, a SQL relation with a
variant_id column such as a backtick-quoted table id. With variants_relation and a coverage backend,
one query finds the variants that have records and reads those records. Only those variants are then
annotated. Without coverage, every variant id in the relation is read into Python and annotated.
staging_dataset is the project.dataset that holds scratch tables. Each table expires after
expiration as a backstop. The table is dropped when the block exits, including on error or
cancellation. With expiration=None a table left behind by a crashed process is never removed.
Raises ValueError if source returns a variant it was not asked for, and SchemaCollisionError if it
returns undeclared columns.
decode_link_cursor
¶
decode_link_cursor(
query: VariantGeneLinkQuery,
) -> str | None
Decode and validate query.cursor, returning its source record id.
encode_link_cursor
¶
encode_link_cursor(
query: VariantGeneLinkQuery, last_record_id: str
) -> str
Encode an opaque cursor bound to the query's locus and filters.
get_variant_gene_link_source_registry
¶
get_variant_gene_link_source_registry() -> (
PluginRegistry[type[VariantGeneLinkSource]]
)
Return the registry for altar.variant_gene_link_sources adapters.
normalize_chromosome
¶
Return a canonical chromosome label without a chr prefix.