Source code for knowledgespaces.assessment.workflows

"""Bounded assessment sessions, observed-pattern replay and latent-state simulation.

These workflows compose the documented questioning and update rules; a
stopping threshold is a decision rule, not a calibration or convergence proof.
"""

from __future__ import annotations

import random
from collections.abc import Callable, Iterable, Sequence
from dataclasses import dataclass
from time import perf_counter
from typing import TYPE_CHECKING, Literal

import numpy as np

from knowledgespaces.assessment.adaptive import (
    is_converged,
    select_item_eig,
    select_item_half_split,
    select_item_informative,
)
from knowledgespaces.assessment.blim import BLIM, StatePosterior, _binary_response
from knowledgespaces.assessment.instances import InstancePool, select_instance_eig
from knowledgespaces.assessment.multiplicative import MultiplicativePosterior

if TYPE_CHECKING:
    from knowledgespaces.estimation.blim_em import ResponseMatrix

Posterior = StatePosterior | MultiplicativePosterior
Selection = Literal["eig", "half_split", "informative"]
StopReason = Literal["threshold", "exhausted", "max_questions"]


[docs] @dataclass(frozen=True) class AssessmentStep: """One observation; times are wall-clock seconds, score uses the chosen policy.""" question: str item: str response: bool score: float selection_seconds: float update_seconds: float
[docs] @dataclass(frozen=True) class AssessmentResult: """Final distribution and every administered question, including unfinished runs. ``probability_history`` is empty unless requested; otherwise it includes the initial distribution and every update, ordered as ``posterior.states``. MAP ties select the first state in cardinality/label order. ``converged`` means only that the requested probability threshold was reached. """ posterior: Posterior steps: tuple[AssessmentStep, ...] stop_reason: StopReason probability_history: tuple[np.ndarray, ...] @property def state(self) -> frozenset[str]: return self.posterior.most_likely_state[0] @property def probability(self) -> float: return self.posterior.most_likely_state[1] @property def converged(self) -> bool: return self.stop_reason == "threshold" @property def questions_asked(self) -> int: return len(self.steps) @property def mean_selection_seconds(self) -> float: return sum(step.selection_seconds for step in self.steps) / max(1, len(self.steps)) @property def mean_update_seconds(self) -> float: return sum(step.update_seconds for step in self.steps) / max(1, len(self.steps))
[docs] def run_assessment( posterior: Posterior, ask_fn: Callable[[str], bool], *, selection: Selection = "half_split", threshold: float = 0.85, max_questions: int = 25, repeat_items: bool = False, instances: InstancePool | None = None, rng: random.Random | None = None, track_probabilities: bool = False, max_memory_bytes: int = 512_000_000, ) -> AssessmentResult: """Run a session from any Bayesian or multiplicative starting distribution. ``ask_fn`` receives an item label, or an instance ID when ``instances`` is supplied. Responses must be booleans or numeric 0/1; other values and callback exceptions propagate. No state is changed in the input. EIG requires a Bayesian posterior. Half-split and informative support both update rules. In item mode ties use candidate order unless ``rng`` is supplied. Instance EIG retains :func:`select_instance_eig`'s random near-tie selection; other instance policies use exact score ties. Pass ``random.Random(seed)`` to reproduce every random selection. By default each item is asked at most once. ``repeat_items=True`` calls ``ask_fn`` anew each time; Bayesian updates then assume independent fresh responses conditional on a fixed state. With an instance pool, each instance is used at most once and ``repeat_items`` must be false. Stop at ``max(P) >= threshold``, exhaustion, or ``max_questions``. The threshold has priority, then exhaustion, then the question limit. A zero question limit is permitted. All options are validated even if the initial distribution already reaches the threshold. Requested probability histories have an estimated memory guard, not an exact process-memory bound. Timing excludes the user's response callback. """ _validate_options( posterior, selection, threshold, max_questions, repeat_items, track_probabilities, max_memory_bytes, ) if not callable(ask_fn): raise TypeError("ask_fn must be callable.") if instances is not None: instances.validate_domain(posterior.structure.domain) if repeat_items: raise ValueError("repeat_items cannot be used with an instance pool.") _memory_guard(posterior, 1, max_questions, track_probabilities, max_memory_bytes) steps: list[AssessmentStep] = [] history = [_readonly(posterior.probabilities)] if track_probabilities else [] asked: set[str] = set() items = sorted(posterior.structure.domain) while True: if is_converged(posterior, threshold): reason: StopReason = "threshold" break candidates = [ q for q in items if ( any(i not in asked for i in instances.instances_for(q)) if instances is not None else repeat_items or q not in asked ) ] if not candidates: reason = "exhausted" break if len(steps) >= max_questions: reason = "max_questions" break start = perf_counter() if selection == "eig" and isinstance(posterior, StatePosterior): if instances is not None: instance = select_instance_eig(posterior, instances, asked=asked, rng=rng) item, question, score = instance.item, instance.instance_id, instance.score else: pick = select_item_eig(posterior, candidates, rng=rng) item, question, score = pick.item, pick.item, pick.score else: policy = ( select_item_half_split if selection == "half_split" else select_item_informative ) pick = policy(posterior, candidates, rng=rng) item, score = pick.item, pick.score available = ( [i for i in instances.instances_for(item) if i not in asked] if instances is not None else [item] ) question = rng.choice(available) if rng is not None else available[0] selection_time = perf_counter() - start response = _binary_response(ask_fn(question)) start = perf_counter() posterior = posterior.update(item, response) update_time = perf_counter() - start asked.add(question) steps.append(AssessmentStep(question, item, response, score, selection_time, update_time)) if track_probabilities: history.append(_readonly(posterior.probabilities)) return AssessmentResult(posterior, tuple(steps), reason, tuple(history))
[docs] @dataclass(frozen=True) class AssessmentEvaluation: """Per-row results with weights and explicitly distinguished distance metrics. Replay reports Hamming distance of the diagnosed state to the observed response pattern, the pattern's minimum distance to the structure, and their difference. These are not latent-state accuracy measurements. Simulation instead reports distance to the supplied true latent state. No failed-threshold session is dropped. Weighted means use all rows. """ results: tuple[AssessmentResult, ...] weights: np.ndarray response_distances: np.ndarray | None = None structure_distances: np.ndarray | None = None true_states: tuple[frozenset[str], ...] | None = None state_distances: np.ndarray | None = None @property def mean_questions(self) -> float: return float(np.average([r.questions_asked for r in self.results], weights=self.weights)) @property def convergence_rate(self) -> float: return float(np.average([r.converged for r in self.results], weights=self.weights)) @property def excess_response_distances(self) -> np.ndarray | None: if self.response_distances is None or self.structure_distances is None: return None return _readonly(self.response_distances - self.structure_distances) @property def mean_response_distance(self) -> float | None: return self._mean(self.response_distances) @property def mean_structure_distance(self) -> float | None: return self._mean(self.structure_distances) @property def mean_excess_response_distance(self) -> float | None: return self._mean(self.excess_response_distances) @property def mean_state_distance(self) -> float | None: return self._mean(self.state_distances) @property def exact_state_rate(self) -> float | None: return None if self.state_distances is None else self._mean(self.state_distances == 0) def _mean(self, values: np.ndarray | None) -> float | None: return None if values is None else float(np.average(values, weights=self.weights))
[docs] def evaluate_assessment( posterior: Posterior, data: ResponseMatrix, *, selection: Selection = "half_split", threshold: float = 0.85, max_questions: int = 25, repeat_items: bool = False, rng: random.Random | None = None, track_probabilities: bool = False, max_memory_bytes: int = 512_000_000, ) -> AssessmentEvaluation: """Replay one assessment per stored complete binary response row. Labels are aligned by name; supplied row weights are retained. Aggregated rows are assessed once, so randomized ties are not independent draws for each individual represented by their weight. Expand rows when those independent draws are wanted. Zero-weight rows remain in the result. ``repeat_items=True`` reproduces fixed-pattern repeated querying as in kstMatrix's assessment simulation. Reusing the same observed answer is not fresh independent evidence: report it as a replay experiment, not calibrated Bayesian assessment or true-state accuracy. Default false prevents that reuse. Every row starts from the same supplied prior. """ from knowledgespaces.estimation.blim_em import ResponseMatrix _validate_options( posterior, selection, threshold, max_questions, repeat_items, track_probabilities, max_memory_bytes, ) # ResponseMatrix arrays are mutable; validate their current contents again. validated = ResponseMatrix(list(data.items), np.asarray(data.patterns), data.counts) if set(validated.items) != posterior.structure.domain: raise ValueError("Response labels must match the assessment domain exactly.") _memory_guard( posterior, validated.n_patterns, max_questions, track_probabilities, max_memory_bytes ) weights = np.ones(validated.n_patterns) if validated.counts is None else validated.counts results, distances, minima = [], [], [] for row in validated.patterns: responses = dict(zip(validated.items, map(bool, row), strict=True)) result = run_assessment( posterior, responses.__getitem__, selection=selection, threshold=threshold, max_questions=max_questions, repeat_items=repeat_items, rng=rng, track_probabilities=track_probabilities, max_memory_bytes=max_memory_bytes, ) observed = frozenset(q for q, correct in responses.items() if correct) results.append(result) distances.append(len(observed ^ result.state)) minima.append(min(len(observed ^ state) for state in posterior.states)) return AssessmentEvaluation( tuple(results), _readonly(weights), _readonly(np.asarray(distances)), _readonly(np.asarray(minima)), )
[docs] def simulate_assessment( posterior: Posterior, true_states: Sequence[Iterable[str]], *, response_model: BLIM, selection: Selection = "half_split", threshold: float = 0.85, max_questions: int = 25, repeat_items: bool = False, seed: int | None = None, track_probabilities: bool = False, max_memory_bytes: int = 512_000_000, ) -> AssessmentEvaluation: """Generate fresh conditionally independent responses during each session. Each supplied state defines one simulated person and remains fixed throughout that assessment. Draw a Bernoulli response only when asked, using ``response_model`` slip/guess parameters. Repeated questions get new draws. Generating and assessing models may differ but must share item labels; true states must belong to the generating model, and need not belong to the assessed structure. This permits misspecification studies. ``seed`` initializes separate local NumPy response and Python tie-breaking generators; it does not reproduce R seeds. Results describe this sample, without an automatic calibration or uncertainty claim. """ _validate_options( posterior, selection, threshold, max_questions, repeat_items, track_probabilities, max_memory_bytes, ) if response_model.structure.domain != posterior.structure.domain: raise ValueError("Generating and assessing models must share the item domain.") if not len(true_states): raise ValueError("Supply at least one true state.") _memory_guard(posterior, len(true_states), max_questions, track_probabilities, max_memory_bytes) states = tuple(frozenset(state) for state in true_states) if any(state not in response_model.structure.states for state in states): raise ValueError("Every true state must belong to the generating model.") if seed is not None and (isinstance(seed, bool) or not isinstance(seed, int) or seed < 0): raise ValueError("seed must be a nonnegative integer or None.") responses_rng, selection_rng = np.random.default_rng(seed), random.Random(seed) results = [] for state in states: def ask(item: str, true_state: frozenset[str] = state) -> bool: return bool(responses_rng.random() < response_model.likelihood(item, True, true_state)) results.append( run_assessment( posterior, ask, selection=selection, threshold=threshold, max_questions=max_questions, repeat_items=repeat_items, rng=selection_rng, track_probabilities=track_probabilities, max_memory_bytes=max_memory_bytes, ) ) distances = np.asarray( [len(state ^ result.state) for state, result in zip(states, results, strict=True)] ) return AssessmentEvaluation( tuple(results), _readonly(np.ones(len(states))), true_states=states, state_distances=_readonly(distances), )
def _validate_options( posterior: Posterior, selection: str, threshold: float, max_questions: int, repeat_items: bool, track_probabilities: bool, max_memory_bytes: int, ) -> None: if selection not in ("eig", "half_split", "informative"): raise ValueError("selection must be 'eig', 'half_split', or 'informative'.") if selection == "eig" and not isinstance(posterior, StatePosterior): raise ValueError("EIG requires a Bayesian StatePosterior with a response model.") is_converged(posterior, threshold) # Validate even with no available questions. if isinstance(max_questions, bool) or not isinstance(max_questions, int) or max_questions < 0: raise ValueError("max_questions must be a nonnegative integer.") if not isinstance(repeat_items, bool) or not isinstance(track_probabilities, bool): raise ValueError("repeat_items and track_probabilities must be booleans.") if ( isinstance(max_memory_bytes, bool) or not isinstance(max_memory_bytes, int) or max_memory_bytes <= 0 ): raise ValueError("max_memory_bytes must be a positive integer.") def _memory_guard( posterior: Posterior, rows: int, max_questions: int, track_probabilities: bool, max_memory_bytes: int, ) -> None: snapshots = max_questions + 2 if track_probabilities else 1 estimate = rows * ( posterior.structure.n_states * 8 * snapshots + max_questions * 256 + posterior.structure.n_items * 256 ) if estimate > max_memory_bytes: raise ValueError( "Estimated assessment output exceeds max_memory_bytes; reduce rows, " "max_questions, or probability tracking, or raise the limit." ) def _readonly(values: np.ndarray) -> np.ndarray: array = np.array(values, copy=True) array.flags.writeable = False return array.view()