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
¶
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
¶
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
¶
Unscored variants whose primary row was cached but a declared detail table had no row.
ineligible_count
class-attribute
instance-attribute
¶
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
¶
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
¶
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.