"""G5 역량별 자기보정·미지 사례 전이·합성 subgroup drift 코어.""" from __future__ import annotations import json import math from collections import defaultdict from pathlib import Path from statistics import mean, stdev from typing import Iterable from ..contracts.calibration_transfer import ( ActualTransferAssessment, ActualTransferExecution, CalibrationBlockInput, CalibrationPair, CalibrationTransferBenchmarkPack, CompetencyCalibrationAssessment, ConfidenceInterval, MetacognitivePrescription, SubgroupDriftReport, SyntheticSubgroupResult, TransferAssessment, TransferSuiteInput, TransferTrial, ) MIN_CALIBRATION_PAIRS = 3 MIN_IMPROVEMENT_PAIRS = 4 MIN_TRANSFER_TRIALS = 4 MIN_ACTUAL_TRANSFER_EXECUTIONS = 4 MIN_SUBGROUP_SAMPLES = 3 TRANSFER_TARGET_RATE = 0.85 DRIFT_GAP_THRESHOLD = 0.20 def _bounded_normal_interval(values: list[float]) -> ConfidenceInterval: if len(values) < 2: center = values[0] return ConfidenceInterval( method="normal_95_bounded", lower=center, upper=center ) margin = 1.96 * stdev(values) / math.sqrt(len(values)) center = mean(values) return ConfidenceInterval( method="normal_95_bounded", lower=max(0.0, center - margin), upper=min(1.0, center + margin), ) def _wilson_interval(successes: int, total: int) -> ConfidenceInterval: if total <= 0: raise ValueError("Wilson interval requires observed trials") z = 1.96 proportion = successes / total denominator = 1 + z * z / total center = (proportion + z * z / (2 * total)) / denominator margin = ( z * math.sqrt(proportion * (1 - proportion) / total + z * z / (4 * total * total)) / denominator ) return ConfidenceInterval( method="wilson_95", lower=max(0.0, center - margin), upper=min(1.0, center + margin), ) def assess_calibration( blocks: Iterable[CalibrationBlockInput], ) -> tuple[CompetencyCalibrationAssessment, ...]: """잠긴 예측과 독립 관찰을 역량별로 대조한다. 총점은 만들지 않는다.""" grouped: dict[str, list[CalibrationBlockInput]] = defaultdict(list) seen_blocks: set[str] = set() for block in sorted(blocks, key=lambda item: item.block_sequence): block_id = block.observation.practice_block_id if block_id in seen_blocks: raise ValueError(f"duplicate calibration practice block: {block_id}") seen_blocks.add(block_id) grouped[block.observation.competency_id].append(block) output: list[CompetencyCalibrationAssessment] = [] for competency_id in sorted(grouped): pairs: list[CalibrationPair] = [] excluded: list[str] = [] for block in grouped[competency_id]: observation = block.observation if observation.status == "insufficient_evidence": excluded.append(observation.practice_block_id) continue prediction = block.prediction_history.locked_prediction observed_success = observation.status == "passed" signed_error = prediction.predicted_success_probability - float( observed_success ) pairs.append( CalibrationPair( practice_block_id=observation.practice_block_id, prediction_id=prediction.prediction_id, observation_id=observation.observation_id, predicted_success_probability=( prediction.predicted_success_probability ), observed_success=observed_success, signed_error=signed_error, absolute_error=abs(signed_error), confidence=prediction.confidence, evidence_refs=observation.evidence_refs, ) ) if len(pairs) < MIN_CALIBRATION_PAIRS: output.append( CompetencyCalibrationAssessment( competency_id=competency_id, pair_count=len(pairs), bias="insufficient_evidence", improvement="insufficient_evidence", pairs=tuple(pairs), excluded_block_ids=tuple(excluded), counterevidence=( f"calibration_pairs_below_minimum:{len(pairs)}/{MIN_CALIBRATION_PAIRS}", ), ) ) continue absolute_errors = [item.absolute_error for item in pairs] signed_errors = [item.signed_error for item in pairs] mean_absolute_error = mean(absolute_errors) mean_signed_error = mean(signed_errors) if mean_signed_error > 0.10: bias = "overconfident" elif mean_signed_error < -0.10: bias = "underconfident" else: bias = "aligned" baseline_error = None recent_error = None improvement = "insufficient_evidence" counterevidence: list[str] = [] if len(pairs) >= MIN_IMPROVEMENT_PAIRS: midpoint = len(pairs) // 2 baseline_error = mean(item.absolute_error for item in pairs[:midpoint]) recent_error = mean(item.absolute_error for item in pairs[midpoint:]) improvement = ( "improved" if recent_error <= baseline_error - 0.05 else "not_improved" ) if improvement == "not_improved": counterevidence.append("recent_calibration_error_did_not_decrease") else: counterevidence.append( f"improvement_pairs_below_minimum:{len(pairs)}/{MIN_IMPROVEMENT_PAIRS}" ) output.append( CompetencyCalibrationAssessment( competency_id=competency_id, pair_count=len(pairs), mean_absolute_error=mean_absolute_error, mean_signed_error=mean_signed_error, error_interval=_bounded_normal_interval(absolute_errors), bias=bias, baseline_error=baseline_error, recent_error=recent_error, improvement=improvement, pairs=tuple(pairs), excluded_block_ids=tuple(excluded), counterevidence=tuple(counterevidence), ) ) return tuple(output) def prescribe_metacognitive_practice( assessment: CompetencyCalibrationAssessment, ) -> MetacognitivePrescription: if assessment.bias == "overconfident": return MetacognitivePrescription( competency_id=assessment.competency_id, bias=assessment.bias, practice_mode="counterevidence_forecast", instruction_ko=( "성공을 예측하기 전에 실패할 수 있는 장면 근거 두 가지를 먼저 적고, " "그 근거를 반영해 성공 확률 범위를 다시 잠가라." ), completion_evidence=( "two_counterevidence_refs", "revised_probability_range_before_reveal", ), ) if assessment.bias == "underconfident": return MetacognitivePrescription( competency_id=assessment.competency_id, bias=assessment.bias, practice_mode="evidence_recall", instruction_ko=( "성공한 미지 장면의 행동 근거와 내담자 반응을 각각 하나씩 회상한 뒤, " "그 근거만으로 다음 성공 확률을 잠가라." ), completion_evidence=( "learner_behavior_evidence_ref", "client_response_evidence_ref", ), ) if assessment.bias == "aligned": return MetacognitivePrescription( competency_id=assessment.competency_id, bias=assessment.bias, practice_mode="uncertainty_range", instruction_ko=( "단일 확신값 대신 성공 가능 범위와 그 범위를 넓히는 불확실성 근거를 " "먼저 기록하고 외부평가 공개 전 잠가라." ), completion_evidence=("probability_interval", "uncertainty_evidence_ref"), ) return MetacognitivePrescription( competency_id=assessment.competency_id, bias=assessment.bias, practice_mode="collect_more_evidence", instruction_ko=( "현재는 자기보정 판정에 필요한 독립 관찰이 부족하다. 같은 역량의 새 장면을 " "최소 세 번 수행하고 각 예측을 외부평가 전에 잠가라." ), completion_evidence=( "three_locked_predictions", "three_independent_observations", ), ) def _coverage(trials: list[TransferTrial]) -> dict[str, int]: return { "contexts": len({item.variation.context_variant for item in trials}), "relationship_styles": len( {item.variation.relationship_style for item in trials} ), "difficulty_levels": len({item.variation.difficulty_level for item in trials}), "expression_variants": len( {item.variation.expression_variant for item in trials} ), "scenario_families": len( {item.variation.scenario_family_id for item in trials} ), "synthetic_subgroups": len( {item.variation.synthetic_subgroup for item in trials} ), } def assess_transfer( suite: TransferSuiteInput, ) -> tuple[TransferAssessment, ...]: grouped: dict[str, list[TransferTrial]] = defaultdict(list) for trial in suite.trials: grouped[trial.competency_id].append(trial) output: list[TransferAssessment] = [] training_phrases = set(suite.training_phrase_family_ids) for competency_id in sorted(grouped): trials = grouped[competency_id] observed = [item for item in trials if item.status != "insufficient_evidence"] successes = sum(item.status == "passed" for item in observed) success_rate = successes / len(observed) if observed else None coverage = _coverage(observed) blockers: list[str] = [] if len(observed) < MIN_TRANSFER_TRIALS: blockers.append( f"observed_trials_below_minimum:{len(observed)}/{MIN_TRANSFER_TRIALS}" ) for dimension in ( "contexts", "relationship_styles", "difficulty_levels", "expression_variants", "scenario_families", ): if coverage[dimension] < 2: blockers.append(f"transfer_coverage_missing:{dimension}") reused = sorted( { item.variation.phrase_family_id for item in observed if item.variation.phrase_family_id in training_phrases } ) if reused: blockers.append("memorized_training_phrase_reused") if any(item.uncertainty > 0.5 for item in observed): blockers.append("transfer_trial_uncertainty_above_boundary") eligible = not blockers transfer_verified = bool( eligible and success_rate is not None and success_rate >= TRANSFER_TARGET_RATE ) if eligible and not transfer_verified: blockers.append( f"transfer_success_rate_below_target:{success_rate:.3f}/{TRANSFER_TARGET_RATE:.2f}" ) output.append( TransferAssessment( competency_id=competency_id, trial_count=len(trials), observed_trial_count=len(observed), success_rate=success_rate, success_interval=( _wilson_interval(successes, len(observed)) if observed else None ), coverage=coverage, eligible=eligible, transfer_verified=transfer_verified, blockers=tuple(blockers), evidence_refs=tuple( dict.fromkeys( ref for item in observed for ref in item.evidence_refs ) ), counterevidence=tuple( dict.fromkeys( ref for item in observed for ref in item.counterevidence ) ), ) ) return tuple(output) def _actual_coverage( executions: list[ActualTransferExecution], ) -> dict[str, int]: return { "contexts": len({item.variation.context_variant for item in executions}), "relationship_styles": len( {item.variation.relationship_style for item in executions} ), "difficulty_levels": len( {item.variation.difficulty_level for item in executions} ), "expression_variants": len( {item.variation.expression_variant for item in executions} ), "scenario_families": len( {item.variation.scenario_family_id for item in executions} ), "synthetic_subgroups": len( {item.variation.synthetic_subgroup for item in executions} ), "phrase_families": len( {item.variation.phrase_family_id for item in executions} ), } def assess_actual_transfer_executions( executions: Iterable[ActualTransferExecution], ) -> tuple[ActualTransferAssessment, ...]: """실제 회기 원장을 역량별로 집계한다. 같은 phrase family의 반복은 최신 실행 하나만 독립 표본으로 인정한다. 합성 benchmark 결과와 섞지 않으며, 훈련 phrase 충돌은 보존하되 전이 분류를 막는다. """ grouped: dict[str, list[ActualTransferExecution]] = defaultdict(list) for execution in executions: grouped[execution.competency_id].append(execution) output: list[ActualTransferAssessment] = [] for competency_id in sorted(grouped): all_items = sorted( grouped[competency_id], key=lambda item: (item.created_at, str(item.execution_event_id)), ) latest_by_phrase: dict[str, ActualTransferExecution] = {} for item in all_items: latest_by_phrase[item.variation.phrase_family_id] = item independent = sorted( latest_by_phrase.values(), key=lambda item: (item.created_at, str(item.execution_event_id)), ) observed = [ item for item in independent if item.status != "insufficient_evidence" and not item.training_phrase_collision ] coverage = _actual_coverage(independent) blockers: list[str] = [] if len(independent) < MIN_ACTUAL_TRANSFER_EXECUTIONS: blockers.append( "actual_independent_phrase_families_below_minimum:" f"{len(independent)}/{MIN_ACTUAL_TRANSFER_EXECUTIONS}" ) if len(observed) < MIN_ACTUAL_TRANSFER_EXECUTIONS: blockers.append( "actual_observed_executions_below_minimum:" f"{len(observed)}/{MIN_ACTUAL_TRANSFER_EXECUTIONS}" ) for dimension in ( "contexts", "relationship_styles", "difficulty_levels", "expression_variants", "scenario_families", ): if coverage[dimension] < 2: blockers.append(f"actual_transfer_coverage_missing:{dimension}") if any(item.training_phrase_collision for item in all_items): blockers.append("training_phrase_family_reused") if any(item.uncertainty > 0.5 for item in observed): blockers.append("actual_transfer_uncertainty_above_boundary") successes = sum(item.status == "passed" for item in observed) success_rate = successes / len(observed) if observed else None eligible = not blockers if not eligible: status = "insufficient_evidence" elif success_rate is not None and success_rate >= TRANSFER_TARGET_RATE: status = "verified" else: status = "not_verified" blockers.append( "actual_transfer_success_rate_below_target:" f"{success_rate:.3f}/{TRANSFER_TARGET_RATE:.2f}" ) output.append( ActualTransferAssessment( competency_id=competency_id, execution_count=len(all_items), independent_execution_count=len(independent), observed_execution_count=len(observed), success_rate=success_rate, success_interval=( _wilson_interval(successes, len(observed)) if observed else None ), coverage=coverage, phrase_family_collision_count=len(all_items) - len(independent), eligible=eligible, actual_transfer_status=status, blockers=tuple(blockers), source_execution_event_ids=tuple( item.execution_event_id for item in all_items ), evidence_turn_ids=tuple( dict.fromkeys( turn_id for item in observed for turn_id in item.evidence_turn_ids ) ), ) ) return tuple(output) def assess_synthetic_subgroup_drift( suite: TransferSuiteInput, ) -> tuple[SubgroupDriftReport, ...]: by_competency: dict[str, list[TransferTrial]] = defaultdict(list) for trial in suite.trials: by_competency[trial.competency_id].append(trial) reports: list[SubgroupDriftReport] = [] for competency_id in sorted(by_competency): by_group: dict[str, list[TransferTrial]] = defaultdict(list) for trial in by_competency[competency_id]: if trial.status != "insufficient_evidence": by_group[trial.variation.synthetic_subgroup].append(trial) results: list[SyntheticSubgroupResult] = [] eligible_rates: dict[str, float] = {} for subgroup in sorted(by_group): items = by_group[subgroup] successes = sum(item.status == "passed" for item in items) rate = successes / len(items) if items else None results.append( SyntheticSubgroupResult( subgroup=subgroup, observed_count=len(items), success_rate=rate, interval=( _wilson_interval(successes, len(items)) if items else None ), ) ) if len(items) >= MIN_SUBGROUP_SAMPLES and rate is not None: eligible_rates[subgroup] = rate if len(eligible_rates) < 2: reports.append( SubgroupDriftReport( competency_id=competency_id, status="insufficient_evidence", compared_subgroups=tuple(sorted(eligible_rates)), subgroup_results=tuple(results), threshold=DRIFT_GAP_THRESHOLD, notice_ko=( "교육용 합성 subgroup별 관측이 충분하지 않아 평가 드리프트를 " "판정하지 않는다. 실제 인구집단 성능 주장이 아니다." ), ) ) continue gap = max(eligible_rates.values()) - min(eligible_rates.values()) reports.append( SubgroupDriftReport( competency_id=competency_id, status=("drift_flagged" if gap > DRIFT_GAP_THRESHOLD else "stable"), max_rate_gap=gap, compared_subgroups=tuple(sorted(eligible_rates)), subgroup_results=tuple(results), threshold=DRIFT_GAP_THRESHOLD, notice_ko=( "이 차이는 교육용 합성 시나리오의 평가 민감도 신호이며 실제 인구집단의 " "능력·위험·임상 결과 차이를 뜻하지 않는다." ), ) ) return tuple(reports) def load_calibration_transfer_benchmark( path: str | Path, ) -> CalibrationTransferBenchmarkPack: return CalibrationTransferBenchmarkPack.model_validate_json( Path(path).read_text(encoding="utf-8") ) def evaluate_calibration_transfer_benchmark( pack: CalibrationTransferBenchmarkPack, ) -> dict[str, object]: case_results: list[dict[str, object]] = [] expected_count = 0 correct_count = 0 contamination_rejections = 0 memorized_false_verifications = 0 for case in pack.cases: calibration = { item.competency_id: item for item in assess_calibration(case.calibration_blocks) } transfer = ( {item.competency_id: item for item in assess_transfer(case.transfer_suite)} if case.transfer_suite else {} ) drift = ( { item.competency_id: item for item in assess_synthetic_subgroup_drift(case.transfer_suite) } if case.transfer_suite else {} ) checks: list[bool] = [] for competency_id, expected in case.expected.calibration_improved.items(): actual = calibration[competency_id].improvement == "improved" checks.append(actual == expected) for competency_id, expected in case.expected.transfer_verified.items(): actual = transfer[competency_id].transfer_verified checks.append(actual == expected) for competency_id, expected in case.expected.drift_status.items(): actual = drift[competency_id].status checks.append(actual == expected) expected_count += len(checks) correct_count += sum(checks) if "post_reveal_contamination" in case.tags: contamination_rejections += 1 if "memorized_phrase_transfer" in case.tags and any( item.transfer_verified for item in transfer.values() ): memorized_false_verifications += 1 case_results.append( { "case_id": case.case_id, "checks": checks, "calibration": { key: value.model_dump(mode="json") for key, value in calibration.items() }, "transfer": { key: value.model_dump(mode="json") for key, value in transfer.items() }, "drift": { key: value.model_dump(mode="json") for key, value in drift.items() }, } ) return { "schema_version": "vignette.calibration-transfer-benchmark-report.v1", "benchmark_version": pack.version, "data_classification": pack.data_classification, "clinical_claim_allowed": pack.clinical_claim_allowed, "expectation_accuracy": ( correct_count / expected_count if expected_count else 1.0 ), "post_reveal_contamination_rejections": contamination_rejections, "memorized_phrase_false_verifications": memorized_false_verifications, "cases": case_results, } def render_calibration_transfer_benchmark_report(report: dict[str, object]) -> str: return json.dumps(report, ensure_ascii=False, indent=2, sort_keys=True) __all__ = [ "DRIFT_GAP_THRESHOLD", "MIN_CALIBRATION_PAIRS", "MIN_IMPROVEMENT_PAIRS", "MIN_SUBGROUP_SAMPLES", "MIN_TRANSFER_TRIALS", "MIN_ACTUAL_TRANSFER_EXECUTIONS", "TRANSFER_TARGET_RATE", "assess_calibration", "assess_actual_transfer_executions", "assess_synthetic_subgroup_drift", "assess_transfer", "evaluate_calibration_transfer_benchmark", "load_calibration_transfer_benchmark", "prescribe_metacognitive_practice", "render_calibration_transfer_benchmark_report", ]