Skip to content

Scoring engine

Model users and binding authors should import the common execution, decoding, validation, and persistence path from altar.scoring.

scoring

Stable execution, decoding, validation, and routing API for model scoring runs.

INELIGIBLE_VARIANT_COLUMNS module-attribute

INELIGIBLE_VARIANT_COLUMNS = (
    "chr",
    "pos",
    "ref",
    "alt",
    "variant_id",
    "variant_class",
    "reason",
)

Header of the ineligible-variant report staged by :func:prepare_scoring_batches.

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.

DetailResultRouter

DetailResultRouter(store: DetailStore)

Validate and stream one named detail result into a generic detail store.

DetailRunResult dataclass

DetailRunResult(
    name: str, added: int, skipped: int, batches: int
)

Write counts for one named detail table.

PrimaryResultRouter

PrimaryResultRouter(store: ScoreStore)

Validate and stream primary scalar batches into one generic score store.

route async

route(
    batches: AsyncIterable[ResultBatch],
    *,
    plugin: ModelPlugin,
    plugin_identity: PluginRunIdentity,
    model_id: str,
    model_name: str,
    genome_build: str,
) -> tuple[int, int, int]

Return (added, skipped, batches) after bounded validation and persistence.

ResultValidationError

Bases: ValueError

A decoded row violates the binding's versioned primary result schema.

ScoringRunResult dataclass

ScoringRunResult(
    plan: ScoringPlan,
    plugin_identity: PluginRunIdentity,
    added: int,
    skipped: int,
    batches: int,
    output_uris: tuple[str, ...] = (),
    details: tuple[DetailRunResult, ...] = (),
)

Durable identities and write counts produced by one engine run.

PreparedScoringBatch dataclass

PreparedScoringBatch(
    request: ScoringRequest,
    variant_input: Transfer,
    variant_count: int,
)

One dense missing-only input and the request that consumes it.

ScoringPreparation dataclass

ScoringPreparation(
    batches: tuple[PreparedScoringBatch, ...],
    plugin_identity: PluginRunIdentity,
    candidate_count: int,
    duplicate_count: int,
    cache_hit_count: int,
    unscored_count: int,
    reuse_policy: ScoreReusePolicy = ScoreReusePolicy(),
    incomplete_detail_count: int = 0,
    ineligible_count: int = 0,
    ineligible_variants: Transfer | None = None,
)

Summary and executable units produced from one complete candidate set.

incomplete_detail_count class-attribute instance-attribute

incomplete_detail_count: int = 0

Unscored variants whose primary row was cached but a declared detail table had no row.

ineligible_count class-attribute instance-attribute

ineligible_count: int = 0

Unique candidates the manifest's variant_eligibility excludes; never looked up or scheduled.

ineligible_variants class-attribute instance-attribute

ineligible_variants: Transfer | None = None

Staged report of the ineligible candidates, or None when there are none.

A headered TSV with :data:INELIGIBLE_VARIANT_COLUMNS, in candidate order. reason is an :class:~altar.models.IneligibilityReason value and variant_class a :class:~altar.models.VariantClass value. The report always describes the latest preparation: when there are no ineligible candidates, a report left at the same location by an earlier preparation is deleted.

ScoringPreparationError

Bases: ValueError

The candidate set or prepared-batch contract is invalid.

DelimitedDetailResultCodec dataclass

DelimitedDetailResultCodec(
    schema: DetailSchema,
    source_names: Mapping[str, str] = dict(),
    allowed_extra_fields: tuple[str, ...] = (),
    identity_field: str = "variant_id",
    delimiter: str = "\t",
    null_tokens: tuple[str, ...] = ("", "NA", "NaN", "nan"),
)

Strict decoder for one headered, delimited named-detail result.

This is the detail-table counterpart to :class:DelimitedResultCodec. It deliberately shares the scalar parsing rules while returning batches tagged with the manifest's DetailSchema.name.

decode async

decode(
    storage: Storage,
    output_uris: Sequence[str],
    *,
    batch_size: int = 10000,
) -> AsyncIterator[DetailResultBatch]

Stream schema-shaped detail rows from every selected file.

DelimitedResultCodec dataclass

DelimitedResultCodec(
    schema: ResultSchema,
    source_names: Mapping[str, str] = dict(),
    row_identity_fields: Mapping[str, str] = dict(),
    allowed_extra_fields: tuple[str, ...] = (),
    identity_field: str = "variant_id",
    delimiter: str = "\t",
    null_tokens: tuple[str, ...] = ("", "NA", "NaN", "nan"),
    optional_extra_fields: tuple[str, ...] = (),
)

Strict decoder for a headered delimited primary-result file.

source_names maps each declared result name to its runtime-file spelling. row_identity_fields maps transport columns such as chr/pos/ref/alt to scalar dtypes and preserves them for storage adapters that need decomposed loci. allowed_extra_fields names required runtime columns that are intentionally ignored; optional_extra_fields names ignored diagnostics accepted when present. Every other undeclared column is rejected. This makes both preservation and projection explicit binding decisions instead of silently dropping new runtime output.

decode async

decode(
    storage: Storage,
    output_uris: Sequence[str],
    *,
    batch_size: int = 10000,
) -> AsyncIterator[ResultBatch]

Stream every collected file in order and validate its complete input schema.

DetailResultBatch dataclass

DetailResultBatch(
    name: str, rows: tuple[dict[str, object], ...]
)

Rows for one binding-declared named detail table.

DetailResultCodec

Bases: Protocol

Decode collected runtime outputs for one manifest-declared detail table.

decode

decode(
    storage: Storage,
    output_uris: Sequence[str],
    *,
    batch_size: int = 10000,
) -> AsyncIterator[DetailResultBatch]

Yield bounded batches bearing the declared detail-table name.

InlineScoringResult dataclass

InlineScoringResult(
    primary_rows: tuple[dict[str, object], ...],
    details: tuple[DetailResultBatch, ...] = (),
)

Primary rows plus lossless named details returned by an inline binding.

ResultBatch dataclass

ResultBatch(rows: tuple[dict[str, object], ...])

One bounded batch of decoded primary result rows.

ResultDecodeError

Bases: ValueError

A runtime output cannot be decoded as the binding's declared result schema.

ScoringResultCodec

Bases: Protocol

Decode collected runtime outputs into bounded, schema-shaped primary batches.

decode

decode(
    storage: Storage,
    output_uris: Sequence[str],
    *,
    batch_size: int = 10000,
) -> AsyncIterator[ResultBatch]

Yield decoded batches without loading a complete output into memory.

run_scoring async

run_scoring(
    plugin: ModelPlugin,
    request: ScoringRequest,
    *,
    model_name: str,
    score_store: ScoreStore,
    detail_store: DetailStore | None = None,
    backend: ExecutionBackend | None = None,
    storage: Storage | None = None,
    runtime: RuntimeContext | None = None,
    input_transfer: TransferMapper | None = None,
    output_transfer: TransferMapper | None = None,
    spec_mapper: SpecMapper | None = None,
    ledger: TaskLedger | None = None,
    poll_interval_s: float = 1.0,
    max_ticks: int = 3600,
    batch_size: int = 10000,
) -> ScoringRunResult

Build, execute, decode, validate, and persist one model scoring request.

Inline plans need only a plugin, request, and store. Container plans additionally require an execution backend and storage adapter. Transfer mappers translate binding-declared logical transfers into a deployment's concrete input/output URIs without teaching the binding about that deployment. A supplied task ledger must be scoped to this run.

prepare_scoring_batches async

prepare_scoring_batches(
    plugin: ModelPlugin,
    request_template: ScoringRequest,
    *,
    candidate_variants: Transfer,
    score_store: ScoreStore,
    storage: Storage,
    destination_transfer: TransferMapper,
    max_variants_per_batch: int,
    cache_query_size: int = 10000,
    reuse_policy: ScoreReusePolicy = ScoreReusePolicy(),
    detail_store: DetailStore | None = None,
) -> ScoringPreparation

Stage dense, cache-miss-only TSV batches and construct their requests.

request_template identifies the model, genome, output selection, and logical layout, but is never executed. It must not already identify a batch. The function derives its expected scientific identity from a plan, verifies any stored scores are compatible, and creates final :class:ScoringRequest values only after the complete global miss set is known.

The candidate input is a canonical headerless TSV with chr, pos, ref, alt[, variant_id] columns. Its first occurrences are retained in input order and duplicate canonical loci are removed on disk. Candidates that plugin.manifest.variant_eligibility excludes are then set aside before any cache lookup: they are never scheduled, and they are staged through destination_transfer as a report at the layout's ineligible_variants_file and returned as ineligible_variants. Cache queries and batch files are bounded by cache_query_size and max_variants_per_batch respectively; the complete cohort and miss set are held in a temporary SQLite relation rather than process memory.

A binding that declares named detail tables requires the detail_store that :func:~altar.scoring.run_scoring will write to. A variant is then a hit only if its primary row is cached for the current run and every declared detail table, stored under a compatible scientific identity, has at least one row for it; otherwise it is prepared again and the engine's key-level deduplication skips rows that already exist. A variant that legitimately produced zero rows for a table therefore remains a miss on every preparation: that costs compute but never loses rich output. Score reuse is primary-only because detail tables carry no score lineage, so a non-empty reuse_policy is rejected for such bindings rather than pairing reused primary scores with details of a different producer.

Warehouse deployments may implement the same lifecycle with a relation-native export while preserving the returned request semantics, including the eligibility filter.