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