"""G2 longitudinal outcome persistence and cohort-safe read service. The deterministic classifier lives in :mod:`outcome_trajectory`. This module only resolves visible ledger evidence, snapshots its provenance, and appends a new immutable revision when evidence changes or recomputation is requested. """ from __future__ import annotations import hashlib import json from collections.abc import Mapping, Sequence from datetime import datetime from typing import Any from uuid import UUID import asyncpg from .. import db from ..contracts.outcome_trajectory import ( OUTCOME_AXES, LongitudinalOutcomeInput, ObservedSessionOutcome, OutcomeAxis, OutcomeAxisObservation, RelationshipEventType, RelationshipMemoryEvent, SafetySignalReference, SyntheticExpectedArc, ) from ..deps import Principal, Role from .outcome_trajectory import build_role_safe_read_model DEFAULT_EXPECTED_ARC_ID = "oas-g2-arc-001" OUTCOME_CHECKIN_INSTRUMENT_ID = "vignette-session-outcome-checkin" OUTCOME_CHECKIN_INSTRUMENT_VERSION = "1.0.0" NON_CLINICAL_NOTICE_KO = ( "이 궤적은 교육용 합성 기대분포와 시뮬레이션 근거를 비교한 학습 피드백이며, " "실제 임상 규준·진단·치료 효과 또는 예후 판단이 아니다." ) class OutcomeTrajectoryNotFoundError(LookupError): pass class OutcomeTrajectoryStateError(ValueError): pass class OutcomeTrajectoryConflictError(RuntimeError): pass def _value(row: Mapping[str, Any], key: str, default: Any = None) -> Any: try: return row[key] except (KeyError, TypeError): return default def _human_view(principal: Principal) -> str: return "counselor" if principal.role == Role.LEARNER else "supervisor" def _db_computed_role(principal: Principal) -> str: return "instructor" if principal.role == Role.TEACHER else principal.role.value def _submission_hash( *, submission_id: UUID, scores: Mapping[str, float], confidences: Mapping[str, float], evidence_turn_ids: Sequence[UUID], ) -> str: payload = { "submission_id": str(submission_id), "scores": {axis: float(scores[axis]) for axis in OUTCOME_AXES}, "confidences": { axis: float(confidences[axis]) for axis in OUTCOME_AXES }, "evidence_turn_ids": sorted(str(item) for item in evidence_turn_ids), } canonical = json.dumps( payload, ensure_ascii=False, separators=(",", ":"), sort_keys=True ) return hashlib.sha256(canonical.encode("utf-8")).hexdigest() def _existing_submission_measurement_ids( rows: Sequence[Mapping[str, Any]], *, submission_hash: str ) -> list[UUID] | None: if not rows: return None by_axis = {str(_value(row, "dimension")): row for row in rows} if len(rows) != len(OUTCOME_AXES) or set(by_axis) != set(OUTCOME_AXES): raise OutcomeTrajectoryConflictError( "outcome observation submission is incomplete in the ledger" ) hashes = { str((_value(row, "metadata", {}) or {}).get("submission_hash", "")) for row in rows } if hashes != {submission_hash}: raise OutcomeTrajectoryConflictError( "submission_id was already used with different outcome observations" ) return [UUID(str(_value(by_axis[axis], "measurement_id"))) for axis in OUTCOME_AXES] def _missing_observation( *, session_no: int, axis: OutcomeAxis, reason: str, measurement: Mapping[str, Any] | None = None, ) -> tuple[OutcomeAxisObservation, dict[str, Any]]: source_kind = _value(measurement or {}, "source_kind", "observed_runtime") perspective = _value(measurement or {}, "perspective", "runtime_observation") instrument_id = _value( measurement or {}, "instrument_id", "vignette-outcome-evidence-gap" ) instrument_version = _value(measurement or {}, "instrument_version", "1.0.0") model_run_id = _value(measurement or {}, "model_run_id") measurement_id = _value(measurement or {}, "measurement_id") status = "error" if reason.startswith("measurement_error:") else "missing" observation = OutcomeAxisObservation( axis=axis, status=status, value=None, confidence=None, source_kind=source_kind, perspective=perspective, instrument_id=instrument_id, instrument_version=instrument_version, model_run_id=model_run_id, evidence_refs=(), missing_reason=reason, ) snapshot = { "measurement_id": measurement_id, "session_no": session_no, "axis": axis, "status": status, "value": None, "raw_value": None, "scale_min": _value(measurement or {}, "scale_min"), "scale_max": _value(measurement or {}, "scale_max"), "confidence": None, "source_kind": source_kind, "perspective": perspective, "instrument_id": instrument_id, "instrument_version": instrument_version, "model_run_id": model_run_id, "evidence_turn_ids": (), "evidence_refs": (), "missing_reason": reason, "source_created_at": _value(measurement or {}, "created_at"), } return observation, snapshot def _observation_from_measurement( *, session_no: int, axis: OutcomeAxis, measurement: Mapping[str, Any] | None, ) -> tuple[OutcomeAxisObservation, dict[str, Any]]: if measurement is None: return _missing_observation( session_no=session_no, axis=axis, reason="measurement_not_collected", ) ledger_status = str(_value(measurement, "status")) error_code = _value(measurement, "error_code") if ledger_status in {"error", "rejected"}: return _missing_observation( session_no=session_no, axis=axis, reason=f"measurement_error:{error_code or ledger_status}", measurement=measurement, ) if ledger_status != "ready": return _missing_observation( session_no=session_no, axis=axis, reason=f"measurement_{ledger_status}", measurement=measurement, ) raw_value = _value(measurement, "value") scale_min = _value(measurement, "scale_min") scale_max = _value(measurement, "scale_max") confidence = _value(measurement, "confidence") evidence_ids = tuple(_value(measurement, "evidence_turn_ids", ()) or ()) if raw_value is None or scale_min is None or scale_max is None or scale_max <= scale_min: return _missing_observation( session_no=session_no, axis=axis, reason="invalid_measurement_scale", measurement=measurement, ) if confidence is None: return _missing_observation( session_no=session_no, axis=axis, reason="measurement_confidence_missing", measurement=measurement, ) source_kind = str(_value(measurement, "source_kind")) if not evidence_ids and source_kind in {"model_inferred", "agent_reported"}: return _missing_observation( session_no=session_no, axis=axis, reason="measurement_evidence_missing", measurement=measurement, ) normalized = (float(raw_value) - float(scale_min)) / ( float(scale_max) - float(scale_min) ) normalized = round(max(0.0, min(1.0, normalized)), 6) evidence_refs = ( tuple(str(item) for item in evidence_ids) if evidence_ids else (f"measurement:{_value(measurement, 'measurement_id')}",) ) observation = OutcomeAxisObservation( axis=axis, status="observed", value=normalized, confidence=float(confidence), source_kind=_value(measurement, "source_kind"), perspective=_value(measurement, "perspective"), instrument_id=_value(measurement, "instrument_id"), instrument_version=_value(measurement, "instrument_version"), model_run_id=_value(measurement, "model_run_id"), evidence_refs=evidence_refs, ) snapshot = { "measurement_id": _value(measurement, "measurement_id"), "session_no": session_no, "axis": axis, "status": "observed", "value": normalized, "raw_value": float(raw_value), "scale_min": float(scale_min), "scale_max": float(scale_max), "confidence": float(confidence), "source_kind": observation.source_kind, "perspective": observation.perspective, "instrument_id": observation.instrument_id, "instrument_version": observation.instrument_version, "model_run_id": observation.model_run_id, "evidence_turn_ids": evidence_ids, "evidence_refs": evidence_refs, "missing_reason": None, "source_created_at": _value(measurement, "created_at"), } return observation, snapshot def _evidence_fingerprint( *, expected_arc_hash: str, snapshots: Sequence[Mapping[str, Any]] ) -> str: payload = { "expected_arc_hash": expected_arc_hash, "observations": [ { key: ( value.isoformat() if isinstance(value, datetime) else str(value) if isinstance(value, UUID) else [str(item) for item in value] if isinstance(value, (tuple, list)) else value ) for key, value in sorted(snapshot.items()) } for snapshot in snapshots ], } canonical = json.dumps( payload, ensure_ascii=False, separators=(",", ":"), sort_keys=True ) return hashlib.sha256(canonical.encode("utf-8")).hexdigest() def _risk_level(ko_risk_level: int | None) -> str: if ko_risk_level is None or ko_risk_level <= 1: return "low" if ko_risk_level == 2: return "moderate" if ko_risk_level == 3: return "high" return "imminent" def _safety_reference(row: Mapping[str, Any]) -> SafetySignalReference: evidence = ( (f"turn:{_value(row, 'turn_id')}",) if _value(row, "turn_id") else (f"safety-event:{_value(row, 'safety_event_id')}",) ) return SafetySignalReference( safety_event_id=str(_value(row, "safety_event_id")), session_no=int(_value(row, "session_no")), risk_level=_risk_level(_value(row, "ko_risk_level")), escalated=bool(_value(row, "escalated")), evidence_refs=evidence, ) def _relationship_event_for_view( row: Mapping[str, Any], *, view: str ) -> RelationshipMemoryEvent: return RelationshipMemoryEvent( event_id=str(_value(row, "memory_event_id")), session_no=int(_value(row, "session_no")), event_type=_value(row, "event_type"), visible_to=(view,), summaries={view: str(_value(row, "summary"))}, evidence_refs=tuple( str(item) for item in (_value(row, "evidence_turn_ids", ()) or ()) ), resolved_by_event_id=( str(_value(row, "resolved_by_event_id")) if _value(row, "resolved_by_event_id") else None ), ) async def _load_visible_case( conn: asyncpg.Connection, session_id: UUID ) -> tuple[Mapping[str, Any], list[Mapping[str, Any]]]: anchor = await conn.fetchrow( """ SELECT id, case_id, learner_id, session_no FROM app.sessions WHERE id = $1 """, session_id, ) if anchor is None: raise OutcomeTrajectoryNotFoundError("session not found or not visible") if _value(anchor, "case_id") is None: raise OutcomeTrajectoryStateError( "longitudinal outcome requires a case_id across sessions" ) anchor_no = _value(anchor, "session_no") if anchor_no is None or not 1 <= int(anchor_no) <= 5: raise OutcomeTrajectoryStateError( "longitudinal outcome currently covers educational sessions 1..5" ) rows = list( await conn.fetch( """ SELECT id, case_id, learner_id, session_no, started_at, ended_at FROM app.sessions WHERE case_id = $1 AND learner_id = $2 AND session_no BETWEEN 1 AND 5 ORDER BY session_no, started_at, id """, _value(anchor, "case_id"), _value(anchor, "learner_id"), ) ) numbers = [int(_value(row, "session_no")) for row in rows] if len(set(numbers)) != len(numbers): raise OutcomeTrajectoryConflictError( "case contains duplicate session numbers in the 1..5 trajectory" ) if numbers != list(range(1, len(rows) + 1)): raise OutcomeTrajectoryStateError( "outcome sessions must be contiguous and ordered from session 1" ) return anchor, rows async def _load_owned_ended_session( conn: asyncpg.Connection, *, principal: Principal, session_id: UUID, ) -> Mapping[str, Any]: if principal.role != Role.LEARNER: raise OutcomeTrajectoryStateError( "outcome observations can only be submitted by a learner" ) row = await conn.fetchrow( """ SELECT id, case_id, learner_id, session_no, ended_at FROM app.sessions WHERE id = $1 AND learner_id = $2 """, session_id, UUID(principal.user_id), ) if row is None: raise OutcomeTrajectoryNotFoundError("session not found or not owned by learner") if _value(row, "ended_at") is None: raise OutcomeTrajectoryStateError( "outcome observations require an ended session" ) if _value(row, "case_id") is None: raise OutcomeTrajectoryStateError( "outcome observations require a longitudinal case_id" ) return row async def _validate_evidence_turns( conn: asyncpg.Connection, *, session_id: UUID, evidence_turn_ids: Sequence[UUID], ) -> None: if not evidence_turn_ids: return visible_count = await conn.fetchval( """ SELECT count(DISTINCT id) FROM app.turns WHERE session_id = $1 AND id = ANY($2::uuid[]) """, session_id, list(evidence_turn_ids), ) if int(visible_count or 0) != len(evidence_turn_ids): raise OutcomeTrajectoryStateError( "evidence_turn_ids must all belong to the requested session" ) async def _load_latest_measurements( conn: asyncpg.Connection, session_ids: Sequence[UUID] ) -> dict[tuple[UUID, str], Mapping[str, Any]]: rows = await conn.fetch( """ WITH current_leaf AS ( SELECT m.*, row_number() OVER ( PARTITION BY m.session_id, m.dimension ORDER BY m.created_at DESC, m.measurement_id DESC ) AS recency FROM app.measurement_event m WHERE m.session_id = ANY($1::uuid[]) AND m.construct = 'session_outcome' AND m.dimension = ANY($2::text[]) AND NOT EXISTS ( SELECT 1 FROM app.measurement_event child WHERE child.supersedes_id = m.measurement_id ) ) SELECT * FROM current_leaf WHERE recency = 1 """, list(session_ids), list(OUTCOME_AXES), ) return { (UUID(str(_value(row, "session_id"))), str(_value(row, "dimension"))): row for row in rows } async def _load_safety( conn: asyncpg.Connection, session_ids: Sequence[UUID] ) -> list[Mapping[str, Any]]: return list( await conn.fetch( """ SELECT se.id AS safety_event_id, se.session_id, s.session_no, se.turn_id, se.ko_risk_level, se.escalated, se.created_at FROM app.safety_events se JOIN app.sessions s ON s.id = se.session_id WHERE se.session_id = ANY($1::uuid[]) ORDER BY s.session_no, se.created_at, se.id """, list(session_ids), ) ) async def _load_relationship_memory( conn: asyncpg.Connection, *, case_id: UUID, view: str ) -> list[Mapping[str, Any]]: return list( await conn.fetch( """ SELECT e.memory_event_id, s.session_no, e.event_type, p.summary, e.evidence_turn_ids, CASE WHEN repair_projection.projection_id IS NOT NULL THEN visible_repair.memory_event_id ELSE NULL END AS resolved_by_event_id FROM app.relationship_memory_event e JOIN app.sessions s ON s.id = e.session_id JOIN app.relationship_memory_projection p ON p.memory_event_id = e.memory_event_id AND p.ai_view = $2 LEFT JOIN app.relationship_memory_event visible_repair ON visible_repair.resolves_event_id = e.memory_event_id AND $2 = ANY(visible_repair.visible_to) LEFT JOIN app.relationship_memory_projection repair_projection ON repair_projection.memory_event_id = visible_repair.memory_event_id AND repair_projection.ai_view = $2 WHERE e.case_id = $1 ORDER BY s.session_no, e.created_at, e.memory_event_id """, case_id, view, ) ) async def _load_expected_arc( conn: asyncpg.Connection, ) -> tuple[SyntheticExpectedArc, Mapping[str, Any]]: row = await conn.fetchrow( """ SELECT arc_id, title_ko, data_classification, clinical_claim_allowed, provenance_note, expected_arc, content_hash FROM ds.synthetic_outcome_arc WHERE arc_id = $1 """, DEFAULT_EXPECTED_ARC_ID, ) if row is None: raise OutcomeTrajectoryStateError("synthetic expected arc registry is missing") return SyntheticExpectedArc.model_validate(_value(row, "expected_arc")), row def _assessment_payload(model: Any) -> dict[str, Any]: payload = model.model_dump(mode="json") # Safety remains a separate top-level ledger reference. The deterministic # core accepts it for transport but never uses it as a classification input. for session in payload["sessions"]: session["safety_signals"] = [] return payload async def _insert_revision( conn: asyncpg.Connection, *, principal: Principal, anchor: Mapping[str, Any], fingerprint: str, assessment: Mapping[str, Any], snapshots: Sequence[Mapping[str, Any]], latest: Mapping[str, Any] | None, reason: str, ) -> Mapping[str, Any]: revision_no = int(_value(latest or {}, "revision_no", 0)) + 1 row = await conn.fetchrow( """ INSERT INTO app.outcome_trajectory_revision ( anchor_session_id, case_id, learner_id, expected_arc_id, revision_no, supersedes_revision_id, source_fingerprint, assessment, observation_count, missing_observation_count, recompute_reason, computed_by, computed_role ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8::jsonb, $9, $10, $11, $12, $13 ) RETURNING revision_id, revision_no, supersedes_revision_id, source_fingerprint, recompute_reason, computed_at """, _value(anchor, "id"), _value(anchor, "case_id"), _value(anchor, "learner_id"), DEFAULT_EXPECTED_ARC_ID, revision_no, _value(latest or {}, "revision_id"), fingerprint, dict(assessment), len(snapshots), sum(item["status"] != "observed" for item in snapshots), reason, UUID(principal.user_id), _db_computed_role(principal), ) assert row is not None for item in snapshots: await conn.execute( """ INSERT INTO app.outcome_trajectory_observation ( revision_id, measurement_id, session_id, session_no, axis, status, value, raw_value, scale_min, scale_max, confidence, source_kind, perspective, instrument_id, instrument_version, model_run_id, evidence_turn_ids, missing_reason, source_created_at ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17::uuid[], $18, $19 ) """, _value(row, "revision_id"), item["measurement_id"], item["session_id"], item["session_no"], item["axis"], item["status"], item["value"], item["raw_value"], item["scale_min"], item["scale_max"], item["confidence"], item["source_kind"], item["perspective"], item["instrument_id"], item["instrument_version"], item["model_run_id"], list(item["evidence_turn_ids"]), item["missing_reason"], item["source_created_at"], ) return row def _response( *, session_id: UUID, revision: Mapping[str, Any], expected_arc_row: Mapping[str, Any], assessment: Mapping[str, Any], snapshots: Sequence[Mapping[str, Any]], safety: Sequence[SafetySignalReference], relationship_memory: Sequence[Any], ) -> dict[str, Any]: expected_arc = dict(_value(expected_arc_row, "expected_arc")) expected_arc["session_count"] = 5 next_questions = list( dict.fromkeys( question for session in assessment.get("sessions", []) for question in session.get("next_check_questions", []) ) ) observations = [] for item in snapshots: observations.append( { "measurement_id": item["measurement_id"], "session_id": item["session_id"], "session_no": item["session_no"], "axis": item["axis"], "status": item["status"], "value": item["value"], "raw_value": item["raw_value"], "scale_min": item["scale_min"], "scale_max": item["scale_max"], "confidence": item["confidence"], "source_kind": item["source_kind"], "perspective": item["perspective"], "instrument_id": item["instrument_id"], "instrument_version": item["instrument_version"], "model_run_id": item["model_run_id"], "evidence_refs": [str(value) for value in item["evidence_refs"]], "missing_reason": item["missing_reason"], "occurred_at": item["source_created_at"], } ) return { "session_id": session_id, "revision_id": _value(revision, "revision_id"), "revision_no": _value(revision, "revision_no"), "supersedes_revision_id": _value(revision, "supersedes_revision_id"), "source_fingerprint": _value(revision, "source_fingerprint"), "recompute_reason": _value(revision, "recompute_reason"), "computed_at": _value(revision, "computed_at"), "notice_ko": NON_CLINICAL_NOTICE_KO, "expected_arc": expected_arc, "assessment": assessment, "next_questions": next_questions, "observations": observations, "safety_signals": [item.model_dump(mode="json") for item in safety], "relationship_memory": [ item.model_dump(mode="json") for item in relationship_memory ], } async def read_outcome_trajectory( *, principal: Principal, session_id: UUID, force_recompute: bool = False, recompute_reason: str | None = None, ) -> dict[str, Any]: """Read or append the visible case trajectory under human RLS context.""" reason = (recompute_reason or "manual_recompute").strip() if force_recompute and not reason: raise OutcomeTrajectoryStateError("recompute_reason must not be blank") async with db.acquire( role=principal.role.value, user_id=principal.user_id, cohort_ids=principal.cohort_ids, ) as conn: anchor, session_rows = await _load_visible_case(conn, session_id) await conn.execute( "SELECT pg_advisory_xact_lock(hashtextextended($1::text, 0))", str(_value(anchor, "case_id")), ) expected_arc, expected_arc_row = await _load_expected_arc(conn) session_ids = [UUID(str(_value(row, "id"))) for row in session_rows] latest_measurements = await _load_latest_measurements(conn, session_ids) safety_rows = await _load_safety(conn, session_ids) view = _human_view(principal) memory_rows = await _load_relationship_memory( conn, case_id=UUID(str(_value(anchor, "case_id"))), view=view ) safety_by_session: dict[int, list[SafetySignalReference]] = {} safety_all = [_safety_reference(row) for row in safety_rows] for signal in safety_all: safety_by_session.setdefault(signal.session_no, []).append(signal) memory_by_session: dict[int, list[RelationshipMemoryEvent]] = {} memory_all = [ _relationship_event_for_view(row, view=view) for row in memory_rows ] for event in memory_all: memory_by_session.setdefault(event.session_no, []).append(event) observed_sessions: list[ObservedSessionOutcome] = [] snapshots: list[dict[str, Any]] = [] for session in session_rows: current_id = UUID(str(_value(session, "id"))) session_no = int(_value(session, "session_no")) axes = [] for axis in OUTCOME_AXES: observation, snapshot = _observation_from_measurement( session_no=session_no, axis=axis, measurement=latest_measurements.get((current_id, axis)), ) snapshot["session_id"] = current_id axes.append(observation) snapshots.append(snapshot) observed_sessions.append( ObservedSessionOutcome( session_no=session_no, axes=tuple(axes), safety_signals=tuple(safety_by_session.get(session_no, ())), relationship_events=tuple(memory_by_session.get(session_no, ())), ) ) trajectory = LongitudinalOutcomeInput( expected_arc=expected_arc, sessions=tuple(observed_sessions), ) read_model = build_role_safe_read_model(trajectory, view=view) assessment = _assessment_payload(read_model.assessment) fingerprint = _evidence_fingerprint( expected_arc_hash=str(_value(expected_arc_row, "content_hash")), snapshots=snapshots, ) latest = await conn.fetchrow( """ SELECT revision_id, revision_no, supersedes_revision_id, source_fingerprint, assessment, recompute_reason, computed_at FROM app.outcome_trajectory_revision WHERE case_id = $1 ORDER BY revision_no DESC LIMIT 1 """, _value(anchor, "case_id"), ) if latest is not None and not force_recompute and _value( latest, "source_fingerprint" ) == fingerprint: revision = latest assessment = _value(latest, "assessment") else: if not force_recompute: reason = "initial_computation" if latest is None else "evidence_changed" revision = await _insert_revision( conn, principal=principal, anchor=anchor, fingerprint=fingerprint, assessment=assessment, snapshots=snapshots, latest=latest, reason=reason, ) return _response( session_id=session_id, revision=revision, expected_arc_row=expected_arc_row, assessment=assessment, snapshots=snapshots, safety=safety_all, relationship_memory=read_model.relationship_memory, ) async def submit_outcome_observations( *, principal: Principal, session_id: UUID, submission_id: UUID, scores: Mapping[str, float], confidences: Mapping[str, float], evidence_turn_ids: Sequence[UUID] = (), ) -> dict[str, Any]: """Idempotently append a learner's three-axis post-session check-in.""" if set(scores) != set(OUTCOME_AXES) or set(confidences) != set(OUTCOME_AXES): raise OutcomeTrajectoryStateError( "outcome observations require all three outcome axes" ) if len(set(evidence_turn_ids)) != len(evidence_turn_ids): raise OutcomeTrajectoryStateError("evidence_turn_ids must be unique") content_hash = _submission_hash( submission_id=submission_id, scores=scores, confidences=confidences, evidence_turn_ids=evidence_turn_ids, ) async with db.acquire( role=principal.role.value, user_id=principal.user_id, cohort_ids=principal.cohort_ids, ) as conn: await _load_owned_ended_session( conn, principal=principal, session_id=session_id, ) await conn.execute( "SELECT pg_advisory_xact_lock(hashtextextended($1::text, 0))", f"outcome-observation:{session_id}", ) existing = list( await conn.fetch( """ SELECT measurement_id, dimension, metadata FROM app.measurement_event WHERE session_id = $1 AND construct = 'session_outcome' AND source_kind = 'learner_reported' AND perspective = 'learner_self_report' AND instrument_id = $2 AND instrument_version = $3 AND metadata->>'submission_id' = $4 ORDER BY dimension, measurement_id """, session_id, OUTCOME_CHECKIN_INSTRUMENT_ID, OUTCOME_CHECKIN_INSTRUMENT_VERSION, str(submission_id), ) ) measurement_ids = _existing_submission_measurement_ids( existing, submission_hash=content_hash, ) if measurement_ids is None: await _validate_evidence_turns( conn, session_id=session_id, evidence_turn_ids=evidence_turn_ids, ) inserted: dict[str, UUID] = {} for axis in OUTCOME_AXES: prior_id = await conn.fetchval( """ SELECT m.measurement_id FROM app.measurement_event m WHERE m.session_id = $1 AND m.construct = 'session_outcome' AND m.dimension = $2 AND m.source_kind = 'learner_reported' AND m.perspective = 'learner_self_report' AND m.instrument_id = $3 AND m.instrument_version = $4 AND NOT EXISTS ( SELECT 1 FROM app.measurement_event child WHERE child.supersedes_id = m.measurement_id ) ORDER BY m.created_at DESC, m.measurement_id DESC LIMIT 1 """, session_id, axis, OUTCOME_CHECKIN_INSTRUMENT_ID, OUTCOME_CHECKIN_INSTRUMENT_VERSION, ) row = await conn.fetchrow( """ INSERT INTO app.measurement_event ( session_id, supersedes_id, construct, dimension, perspective, source_kind, instrument_id, instrument_version, value, scale_min, scale_max, confidence, status, evidence_turn_ids, visible_to, metadata ) VALUES ( $1, $2, 'session_outcome', $3, 'learner_self_report', 'learner_reported', $4, $5, $6, 0, 1, $7, 'ready', $8::uuid[], ARRAY['counselor','evaluator','supervisor']::text[], $9::jsonb ) RETURNING measurement_id """, session_id, prior_id, axis, OUTCOME_CHECKIN_INSTRUMENT_ID, OUTCOME_CHECKIN_INSTRUMENT_VERSION, float(scores[axis]), float(confidences[axis]), list(evidence_turn_ids), { "submission_id": str(submission_id), "submission_hash": content_hash, "evidence_basis": ( "transcript_turns" if evidence_turn_ids else "learner_self_report_submission" ), "clinical_claim_allowed": False, }, ) assert row is not None inserted[axis] = UUID(str(_value(row, "measurement_id"))) measurement_ids = [inserted[axis] for axis in OUTCOME_AXES] trajectory = await read_outcome_trajectory( principal=principal, session_id=session_id, ) trajectory["submission_id"] = submission_id trajectory["submitted_measurement_ids"] = measurement_ids return trajectory async def append_relationship_memory_event( *, principal: Principal, session_id: UUID, event_type: RelationshipEventType, summaries: Mapping[str, str], evidence_turn_ids: Sequence[UUID], resolves_event_id: UUID | None = None, ) -> UUID: """Append a human-supervisor relationship memory and role projections.""" if principal.role not in {Role.TEACHER, Role.ADMIN}: raise OutcomeTrajectoryStateError( "relationship memory authoring requires teacher or admin role" ) normalized = {key: value.strip() for key, value in summaries.items()} if not normalized or any(not value for value in normalized.values()): raise OutcomeTrajectoryStateError("relationship summaries must not be blank") if len(set(evidence_turn_ids)) != len(evidence_turn_ids) or not evidence_turn_ids: raise OutcomeTrajectoryStateError( "relationship evidence_turn_ids must be non-empty and unique" ) async with db.acquire( role=principal.role.value, user_id=principal.user_id, cohort_ids=principal.cohort_ids, ) as conn: anchor, _ = await _load_visible_case(conn, session_id) await _validate_evidence_turns( conn, session_id=session_id, evidence_turn_ids=evidence_turn_ids, ) try: row = await conn.fetchrow( """ INSERT INTO app.relationship_memory_event ( session_id, case_id, event_type, resolves_event_id, visible_to, evidence_turn_ids, source_kind, created_by ) VALUES ($1, $2, $3, $4, $5::text[], $6::uuid[], 'human_rated', $7) RETURNING memory_event_id """, session_id, _value(anchor, "case_id"), event_type, resolves_event_id, list(normalized), list(evidence_turn_ids), UUID(principal.user_id), ) except (asyncpg.CheckViolationError, asyncpg.ForeignKeyViolationError) as exc: raise OutcomeTrajectoryStateError( "relationship memory resolve/evidence contract was rejected" ) from exc assert row is not None memory_event_id = UUID(str(_value(row, "memory_event_id"))) for view, summary in normalized.items(): await conn.execute( """ INSERT INTO app.relationship_memory_projection ( memory_event_id, ai_view, summary ) VALUES ($1, $2, $3) """, memory_event_id, view, summary, ) return memory_event_id __all__ = [ "DEFAULT_EXPECTED_ARC_ID", "NON_CLINICAL_NOTICE_KO", "OutcomeTrajectoryConflictError", "OutcomeTrajectoryNotFoundError", "OutcomeTrajectoryStateError", "append_relationship_memory_event", "read_outcome_trajectory", "submit_outcome_observations", ]