diff --git a/apps/api/app/access_logging.py b/apps/api/app/access_logging.py new file mode 100644 index 0000000..90c241f --- /dev/null +++ b/apps/api/app/access_logging.py @@ -0,0 +1,60 @@ +"""Uvicorn access log에서 인증 콜백 비밀값을 제거한다. + +Uvicorn은 ASGI scope의 path와 query string을 합쳐 ``uvicorn.access``의 +세 번째 포맷 인자로 기록한다. Google OAuth 콜백에는 일회용 authorization +code와 CSRF state가 query에 있으므로, 애플리케이션 라우터가 안전하게 로깅해도 +기본 access log에서 먼저 노출될 수 있다. +""" + +from __future__ import annotations + +import logging +from typing import Any + + +OAUTH_CALLBACK_PATH = "/auth/callback" + + +def redact_access_log_target(target: Any) -> Any: + """OAuth callback의 query 전체를 제거하고 다른 request target은 보존한다.""" + if not isinstance(target, str): + return target + path, separator, _query = target.partition("?") + if separator and path.rstrip("/") == OAUTH_CALLBACK_PATH: + return path + return target + + +class UvicornAccessLogRedactionFilter(logging.Filter): + """Uvicorn의 구조화된 access-log 인자에서 request target만 정제한다.""" + + def filter(self, record: logging.LogRecord) -> bool: + args = record.args + # Uvicorn 0.30.x access 포맷: + # (client_addr, method, path_with_query, http_version, status_code) + if isinstance(args, tuple) and len(args) >= 3: + target = args[2] + redacted = redact_access_log_target(target) + if redacted != target: + record.args = (*args[:2], redacted, *args[3:]) + return True + + +_ACCESS_LOG_FILTER = UvicornAccessLogRedactionFilter() + + +def install_uvicorn_access_log_redaction() -> None: + """현재 프로세스의 Uvicorn access logger에 필터를 한 번만 설치한다.""" + access_logger = logging.getLogger("uvicorn.access") + if not any( + isinstance(existing, UvicornAccessLogRedactionFilter) + for existing in access_logger.filters + ): + access_logger.addFilter(_ACCESS_LOG_FILTER) + + +__all__ = [ + "UvicornAccessLogRedactionFilter", + "install_uvicorn_access_log_redaction", + "redact_access_log_target", +] diff --git a/apps/api/app/auth_sessions.py b/apps/api/app/auth_sessions.py index 866ecf6..46fc35e 100644 --- a/apps/api/app/auth_sessions.py +++ b/apps/api/app/auth_sessions.py @@ -39,6 +39,7 @@ class SessionUser: consent_at: float | None profile_completed_at: float | None expires_at: float + learner_feedback_enabled: bool = True @dataclass(slots=True) @@ -67,6 +68,7 @@ class ManagedUser: privacy_version: str created_at: float last_seen_at: float + learner_feedback_enabled: bool = True @dataclass(slots=True) @@ -93,6 +95,7 @@ class ManagedUserMemoryInput: privacy_agreed_at: float | None = None terms_version: str | None = None privacy_version: str | None = None + learner_feedback_enabled: bool | None = None reactivate: bool = False @classmethod @@ -125,6 +128,7 @@ class ManagedUserMemoryInput: privacy_agreed_at=user.privacy_agreed_at, terms_version=user.terms_version, privacy_version=user.privacy_version, + learner_feedback_enabled=user.learner_feedback_enabled, reactivate=reactivate, ) @@ -140,6 +144,7 @@ class ManagedUserUpsertInput: external_id: str | None = None affiliation: str | None = None account_status: AccountStatus | None = None + learner_feedback_enabled: bool | None = None reactivate: bool = False @@ -162,6 +167,7 @@ class ManagedUserPatch: complete_onboarding: bool = False terms_version: str | None = None privacy_version: str | None = None + learner_feedback_enabled: bool | None = None class InactiveUserError(Exception): @@ -203,6 +209,21 @@ def _normalize_email(email: str) -> str: return email.strip().lower() +def _email_domain(email: str) -> str: + normalized = _normalize_email(email) + if "@" not in normalized: + return "" + return normalized.rsplit("@", 1)[1] + + +def _allowed_email_domain_set() -> set[str]: + return { + normalized + for value in settings.auth_allowed_email_domains + if (normalized := str(value).strip().lower().lstrip("@")) + } + + def _normalize_external_id(external_id: str | None, email: str) -> str: value = (external_id or "").strip().lower() return value or f"email:{_normalize_email(email)}" @@ -270,7 +291,13 @@ def _initial_account_status( normalized_email = _normalize_email(email) if reactivate or normalized_email in _auto_approved_email_set(): return "approved" - if settings.environment == "dev" and _is_dev_login_external_id(external_id): + # 로컬 E2E 자동 승인은 허용된 내부 도메인에만 적용한다. 관리자 exact-email로 + # 사전등록한 외부 연구참여자의 pending 상태를 dev 로그인이 승격하면 안 된다. + if ( + settings.environment == "dev" + and _is_dev_login_external_id(external_id) + and _email_domain(normalized_email) in _allowed_email_domain_set() + ): return "approved" return _account_status(settings.auth_new_user_default_status) @@ -332,10 +359,11 @@ async def _runtime_tables_ready(conn) -> bool: 'last_seen_at', 'account_status', 'admin_access', + 'learner_feedback_enabled', 'updated_at' ) GROUP BY table_schema, table_name - HAVING count(*) = 19 + HAVING count(*) = 20 ) AS has_user_columns, EXISTS ( SELECT 1 FROM information_schema.columns @@ -362,10 +390,11 @@ async def _runtime_tables_ready(conn) -> bool: 'persona_display_name', 'persona_difficulty', 'prev_rapport_credit', - 'session_goals' + 'session_goals', + 'learner_feedback_enabled' ) GROUP BY table_schema, table_name - HAVING count(*) = 6 + HAVING count(*) = 7 ) AS has_session_columns, EXISTS ( SELECT 1 FROM information_schema.columns @@ -542,10 +571,19 @@ async def ensure_runtime_tables() -> None: """ DO $$ BEGIN - IF to_regclass('app.persona_card') IS NOT NULL THEN + IF to_regclass('app.persona_card') IS NOT NULL + AND NOT EXISTS ( + SELECT 1 FROM information_schema.columns + WHERE table_schema = 'app' + AND table_name = 'persona_card' + AND column_name = 'triggers' + ) THEN ALTER TABLE app.persona_card ADD COLUMN IF NOT EXISTS triggers JSONB NOT NULL DEFAULT '{}'::jsonb; + END IF; + IF to_regclass('app.persona_card') IS NOT NULL + AND to_regclass('app.persona_voice_map') IS NULL THEN CREATE TABLE IF NOT EXISTS app.persona_voice_map ( persona_id UUID NOT NULL, version INT NOT NULL, @@ -564,6 +602,7 @@ async def ensure_runtime_tables() -> None: WHERE table_schema = 'app' AND table_name = 'app_user' AND column_name = 'affiliation' + AND column_default IS DISTINCT FROM ''''::text ) THEN ALTER TABLE app.app_user ALTER COLUMN affiliation SET DEFAULT ''; END IF; @@ -593,6 +632,7 @@ async def ensure_runtime_tables() -> None: ADD COLUMN IF NOT EXISTS last_seen_at TIMESTAMPTZ NOT NULL DEFAULT now(), ADD COLUMN IF NOT EXISTS account_status TEXT NOT NULL DEFAULT 'approved', ADD COLUMN IF NOT EXISTS admin_access BOOLEAN NOT NULL DEFAULT false, + ADD COLUMN IF NOT EXISTS learner_feedback_enabled BOOLEAN NOT NULL DEFAULT true, ADD COLUMN IF NOT EXISTS updated_at TIMESTAMPTZ NOT NULL DEFAULT now() """ ) @@ -926,7 +966,8 @@ async def ensure_runtime_tables() -> None: ADD COLUMN IF NOT EXISTS persona_display_name TEXT, ADD COLUMN IF NOT EXISTS persona_difficulty TEXT, ADD COLUMN IF NOT EXISTS prev_rapport_credit REAL NOT NULL DEFAULT 0.0, - ADD COLUMN IF NOT EXISTS session_goals JSONB NOT NULL DEFAULT '[]'::jsonb + ADD COLUMN IF NOT EXISTS session_goals JSONB NOT NULL DEFAULT '[]'::jsonb, + ADD COLUMN IF NOT EXISTS learner_feedback_enabled BOOLEAN NOT NULL DEFAULT true """ ) await conn.execute( @@ -1108,6 +1149,9 @@ def _managed_user_from_row(row) -> ManagedUser: privacy_version=_row_value(row, "privacy_version", "") or "", created_at=_ts(row["created_at"]), last_seen_at=_ts(row["last_seen_at"]), + learner_feedback_enabled=bool( + _row_value(row, "learner_feedback_enabled", True) + ), ) @@ -1133,8 +1177,14 @@ def _memory_upsert_managed_user(data: ManagedUserMemoryInput) -> ManagedUser: admin_access=data.admin_access if data.admin_access is not None else (current.admin_access if current else False), - account_status=data.account_status - or (current.account_status if current else "approved"), + account_status=( + "suspended" + if current is not None + and current.account_status == "suspended" + and not data.reactivate + else data.account_status + or (current.account_status if current else "approved") + ), cohort_ids=( list(data.cohort_ids) if data.cohort_ids is not None @@ -1209,6 +1259,11 @@ def _memory_upsert_managed_user(data: ManagedUserMemoryInput) -> ManagedUser: ), created_at=current.created_at if current else now, last_seen_at=now, + learner_feedback_enabled=( + data.learner_feedback_enabled + if data.learner_feedback_enabled is not None + else (current.learner_feedback_enabled if current else True) + ), ) _users[uid] = user _email_index[normalized_email] = uid @@ -1232,6 +1287,7 @@ def _sync_managed_user_sessions(user: ManagedUser) -> None: session.cohort_ids = list(user.cohort_ids) session.consent_at = user.consent_at session.profile_completed_at = user.profile_completed_at + session.learner_feedback_enabled = user.learner_feedback_enabled async def upsert_managed_user(data: ManagedUserUpsertInput) -> ManagedUser: @@ -1267,6 +1323,11 @@ async def upsert_managed_user(data: ManagedUserUpsertInput) -> ManagedUser: WHEN $7::boolean IS NULL THEN admin_access ELSE $7::boolean END, + account_status = CASE + WHEN account_status = 'suspended' THEN 'suspended' + WHEN $8::text = 'approved' THEN 'approved' + ELSE account_status + END, last_seen_at = now(), updated_at = now() WHERE user_id = ( @@ -1288,6 +1349,7 @@ async def upsert_managed_user(data: ManagedUserUpsertInput) -> ManagedUser: display_name, role, admin_access, + learner_feedback_enabled, account_status, cohort, affiliation, @@ -1319,6 +1381,7 @@ async def upsert_managed_user(data: ManagedUserUpsertInput) -> ManagedUser: data.affiliation or DEFAULT_AFFILIATION, manual_external_id, desired_admin_access, + desired_account_status, ) if row is not None: user = _managed_user_from_row(row) @@ -1334,13 +1397,17 @@ async def upsert_managed_user(data: ManagedUserUpsertInput) -> ManagedUser: display_name, role, admin_access, + learner_feedback_enabled, cohort, affiliation, account_status, last_seen_at, updated_at ) - VALUES ($1, $2, $3, $4, COALESCE($9::boolean, false), $5, $6, $8, now(), now()) + VALUES ( + $1, $2, $3, $4, COALESCE($9::boolean, false), + COALESCE($10::boolean, true), $5, $6, $8, now(), now() + ) ON CONFLICT (external_id) DO UPDATE SET email = EXCLUDED.email, display_name = COALESCE(NULLIF(EXCLUDED.display_name, ''), app.app_user.display_name), @@ -1349,6 +1416,10 @@ async def upsert_managed_user(data: ManagedUserUpsertInput) -> ManagedUser: WHEN $9::boolean IS NULL THEN app.app_user.admin_access ELSE EXCLUDED.admin_access END, + learner_feedback_enabled = CASE + WHEN $10::boolean IS NULL THEN app.app_user.learner_feedback_enabled + ELSE EXCLUDED.learner_feedback_enabled + END, cohort = COALESCE(EXCLUDED.cohort, app.app_user.cohort), affiliation = COALESCE(NULLIF(EXCLUDED.affiliation, ''), app.app_user.affiliation), is_active = CASE WHEN $7 THEN TRUE ELSE app.app_user.is_active END, @@ -1367,6 +1438,7 @@ async def upsert_managed_user(data: ManagedUserUpsertInput) -> ManagedUser: display_name, role, admin_access, + learner_feedback_enabled, account_status, cohort, affiliation, @@ -1396,6 +1468,7 @@ async def upsert_managed_user(data: ManagedUserUpsertInput) -> ManagedUser: data.reactivate, desired_account_status, desired_admin_access, + data.learner_feedback_enabled, ) if row is None: _inactive_emails.add(normalized_email) @@ -1438,6 +1511,7 @@ async def upsert_managed_user(data: ManagedUserUpsertInput) -> ManagedUser: cohort_ids=data.cohort_ids, user_id=fallback_uid, affiliation=data.affiliation, + learner_feedback_enabled=data.learner_feedback_enabled, reactivate=data.reactivate, ) ) @@ -1455,6 +1529,7 @@ async def get_managed_user(user_id: str) -> ManagedUser | None: display_name, role, admin_access, + learner_feedback_enabled, account_status, cohort, affiliation, @@ -1504,6 +1579,7 @@ async def get_managed_user_by_email(email: str) -> ManagedUser | None: display_name, role, admin_access, + learner_feedback_enabled, account_status, cohort, affiliation, @@ -1555,6 +1631,7 @@ async def list_managed_users() -> tuple[list[ManagedUser], bool]: display_name, role, admin_access, + learner_feedback_enabled, account_status, cohort, affiliation, @@ -1599,6 +1676,7 @@ async def update_managed_user( role = COALESCE($3, role), account_status = COALESCE($18, account_status), admin_access = COALESCE($19, admin_access), + learner_feedback_enabled = COALESCE($20, learner_feedback_enabled), cohort = CASE WHEN $4 THEN $5 ELSE cohort END, affiliation = COALESCE($6, affiliation), legal_name = COALESCE($7, legal_name), @@ -1623,6 +1701,7 @@ async def update_managed_user( display_name, role, admin_access, + learner_feedback_enabled, account_status, cohort, affiliation, @@ -1670,6 +1749,7 @@ async def update_managed_user( else None, patch.account_status, patch.admin_access, + patch.learner_feedback_enabled, ) if row is not None: next_user = _managed_user_from_row(row) @@ -1698,6 +1778,9 @@ async def update_managed_user( account_status=patch.account_status if patch.account_status is not None else current.account_status, + learner_feedback_enabled=patch.learner_feedback_enabled + if patch.learner_feedback_enabled is not None + else current.learner_feedback_enabled, cohort_ids=list(patch.cohort_ids) if patch.cohort_ids is not None else current.cohort_ids, @@ -1934,6 +2017,7 @@ async def create_session( cohort_ids: list[str] | None = None, user_id: str | None = None, external_id: str | None = None, + account_status: AccountStatus | None = None, ) -> tuple[str, SessionUser]: raw_sid = secrets.token_urlsafe(32) normalized_email = _normalize_email(email) @@ -1945,6 +2029,7 @@ async def create_session( cohort_ids=cohort_ids, user_id=user_id, external_id=external_id, + account_status=account_status, reactivate=False, ) ) @@ -1973,6 +2058,7 @@ async def create_session( consent_at=managed.consent_at, profile_completed_at=managed.profile_completed_at, expires_at=expires_at, + learner_feedback_enabled=managed.learner_feedback_enabled, ) sid_hash = _sid_hash(raw_sid) try: @@ -2022,6 +2108,7 @@ async def get_session(raw_sid: str | None) -> SessionUser | None: COALESCE(u.display_name, s.display_name, u.email) AS display_name, u.role, u.admin_access, + u.learner_feedback_enabled, u.account_status, u.cohort, u.consent_at, @@ -2078,6 +2165,9 @@ async def get_session(raw_sid: str | None) -> SessionUser | None: consent_at=_optional_ts(row["consent_at"]), profile_completed_at=_optional_ts(row["profile_completed_at"]), expires_at=_ts(row["expires_at"]), + learner_feedback_enabled=bool( + _row_value(row, "learner_feedback_enabled", True) + ), ) except Exception: require_runtime_fallback_allowed("browser session") @@ -2103,6 +2193,7 @@ async def get_session(raw_sid: str | None) -> SessionUser | None: user.cohort_ids = list(managed.cohort_ids) user.consent_at = managed.consent_at user.profile_completed_at = managed.profile_completed_at + user.learner_feedback_enabled = managed.learner_feedback_enabled return user diff --git a/apps/api/app/config.py b/apps/api/app/config.py index be9673b..ed8eb32 100644 --- a/apps/api/app/config.py +++ b/apps/api/app/config.py @@ -392,6 +392,8 @@ class Settings(BaseSettings): default="https://api-vignette.chanpaca.net/auth/callback", validation_alias="OAUTH_REDIRECT_URI", ) + # 로컬 dev-login과 SAML 조직 정책용 allowlist다. Google OIDC는 provider가 + # 검증한 이메일이면 도메인·사전등록 여부와 무관하게 로그인시킨다. auth_allowed_email_domains: list[str] = Field( default=["hs.ac.kr", "twentyoz.kr"], validation_alias="AUTH_ALLOWED_EMAIL_DOMAINS", diff --git a/apps/api/app/deps.py b/apps/api/app/deps.py index 77b4a94..b4d08cf 100644 --- a/apps/api/app/deps.py +++ b/apps/api/app/deps.py @@ -41,6 +41,7 @@ class Principal: display_name: str = "", consent_at: float | None = None, profile_completed_at: float | None = None, + learner_feedback_enabled: bool = True, ) -> None: self.user_id = user_id self.role = role @@ -52,6 +53,7 @@ class Principal: self.display_name = display_name self.consent_at = consent_at self.profile_completed_at = profile_completed_at + self.learner_feedback_enabled = learner_feedback_enabled def can_access_role(self, role: Role) -> bool: if self.role == role: @@ -76,6 +78,7 @@ class Principal: display_name=self.display_name, consent_at=self.consent_at, profile_completed_at=self.profile_completed_at, + learner_feedback_enabled=self.learner_feedback_enabled, ) @@ -122,6 +125,7 @@ async def get_current_principal( display_name=session.display_name, consent_at=session.consent_at, profile_completed_at=session.profile_completed_at, + learner_feedback_enabled=session.learner_feedback_enabled, ) diff --git a/apps/api/app/engine_client.py b/apps/api/app/engine_client.py index 09b9a14..6f4825a 100644 --- a/apps/api/app/engine_client.py +++ b/apps/api/app/engine_client.py @@ -139,24 +139,38 @@ class EngineClient: async def health(self) -> bool: return bool((await self.health_detail()).get("ok")) - async def health_detail(self) -> dict[str, Any]: + async def _provider_health_detail( + self, + *, + provider: EngineProvider, + model: str | None = None, + reasoning_effort: ReasoningEffort | None = None, + ) -> dict[str, Any]: try: - params: dict[str, str] = {"provider": self.engine_mode} - if self.default_model: - params["model"] = self.default_model - if self.default_reasoning_effort: - params["reasoning_effort"] = self.default_reasoning_effort + params: dict[str, str] = {"provider": provider} + if model: + params["model"] = model + if reasoning_effort: + params["reasoning_effort"] = reasoning_effort r = await self.client.get("/ready", params=params) if r.status_code == 404: live = await self.client.get("/health") return { - "ok": live.status_code == 200, - "detail": "gateway liveness only; readiness endpoint unavailable", + # 게이트웨이 프로세스가 살아 있다는 사실만으로 특정 provider가 + # 응답을 만들 수 있다고 판정하면 회기 화면이 거짓 GREEN이 된다. + "ok": False, + "detail": ( + "게이트웨이는 응답하지만 공급자 준비상태 엔드포인트를 사용할 수 없음" + if live.status_code == 200 + else "게이트웨이와 공급자 준비상태를 확인할 수 없음" + ), "status_code": live.status_code, + "cached": False, } payload: dict[str, Any] = {} try: - payload = r.json() + decoded = r.json() + payload = decoded if isinstance(decoded, dict) else {} except ValueError: payload = {} return { @@ -170,8 +184,62 @@ class EngineClient: "ok": False, "detail": f"engine readiness transport error: {exc}", "status_code": None, + "cached": False, } + async def health_detail(self) -> dict[str, Any]: + """관리자 기본 lane과 실시간 내담자 lane의 준비상태를 함께 반환한다. + + 내담자 턴은 ``live_client_provider``로, 평가·리뷰는 ``engine_mode``로 + 분리될 수 있으므로 기본 공급자 하나만 확인해서는 실제 응답 가능 여부를 알 수 없다. + """ + + default_detail = await self._provider_health_detail( + provider=self.engine_mode, + model=self.default_model, + reasoning_effort=self.default_reasoning_effort, + ) + live_provider = self.live_client_provider or self.engine_mode + if live_provider == self.engine_mode: + live_detail = dict(default_detail) + else: + # 관리자 기본 model/reasoning 값은 engine_mode 소속이므로 별도 client + # provider readiness에는 넘기지 않는다(_payload 계약과 동일). + live_detail = await self._provider_health_detail(provider=live_provider) + + default_ok = bool(default_detail.get("ok")) + live_ok = bool(live_detail.get("ok")) + if not default_ok: + detail = ( + f"기본 공급자 {self.engine_mode}: " + f"{default_detail.get('detail') or 'readiness failed'}" + ) + status_code = default_detail.get("status_code") + elif not live_ok: + detail = ( + f"실시간 내담자 공급자 {live_provider}: " + f"{live_detail.get('detail') or 'readiness failed'}" + ) + status_code = live_detail.get("status_code") + else: + detail = str(live_detail.get("detail") or default_detail.get("detail") or "ready") + status_code = live_detail.get("status_code") + + return { + "ok": default_ok and live_ok, + "detail": detail, + "status_code": status_code, + "cached": bool(default_detail.get("cached")) and bool(live_detail.get("cached")), + "default_engine": { + "provider": self.engine_mode, + **default_detail, + }, + "live_client_engine": { + "provider": live_provider, + **live_detail, + }, + } + async def capabilities( self, *, diff --git a/apps/api/app/main.py b/apps/api/app/main.py index 13303bf..220a86b 100644 --- a/apps/api/app/main.py +++ b/apps/api/app/main.py @@ -16,6 +16,7 @@ from fastapi.middleware.cors import CORSMiddleware from fastapi.staticfiles import StaticFiles from . import __version__ +from .access_logging import install_uvicorn_access_log_redaction from .auth_sessions import ensure_runtime_tables from .config import settings from .db import acquire, close_pool, healthcheck, init_pool @@ -46,6 +47,7 @@ from .routes import multimodal_alliance as multimodal_alliance_routes from .routes import measurements as measurement_routes from .routes import outcome_trajectories as outcome_trajectory_routes from .routes import personas as persona_routes +from .routes import protocols as protocol_routes from .routes import rupture_repairs as rupture_repair_routes from .routes import supervision_research as supervision_research_routes from .routes import sessions as session_routes @@ -54,6 +56,7 @@ from .routes import teacher as teacher_routes from .routes import users as user_routes from .routes import voice as voice_routes from .services.notifications import ensure_notification_tables +from .services.protocol_registry import ensure_protocol_tables from .services import ( alliance_measurement, continuous_improvement_producer, @@ -63,6 +66,7 @@ from .services.voice import voice_service logger = logging.getLogger(__name__) +install_uvicorn_access_log_redaction() @asynccontextmanager @@ -84,6 +88,7 @@ async def lifespan(app: FastAPI): await ensure_runtime_tables() await ensure_review_tables() await ensure_notification_tables() + await ensure_protocol_tables() async with acquire(role="admin") as conn: measurement_schema_ready = await schema_contract_ready( conn, @@ -251,6 +256,7 @@ app.include_router(calibration_transfer_routes.router) app.include_router(continuous_improvement_routes.router) app.include_router(client_diagnostics_routes.router) app.include_router(admin_routes.router) +app.include_router(protocol_routes.router) app.include_router(persona_routes.router) app.include_router(session_routes.router) app.include_router(measurement_routes.router) @@ -281,5 +287,10 @@ async def health() -> dict[str, object]: "db": db_ok, "engine": engine_ok, "engine_detail": engine.get("detail"), - "engine_mode": settings.engine_mode, + "engine_mode": engine_client.engine_mode, + "live_client_provider": ( + engine_client.live_client_provider or engine_client.engine_mode + ), + "default_engine": engine.get("default_engine"), + "live_client_engine": engine.get("live_client_engine"), } diff --git a/apps/api/app/persona_read_model.py b/apps/api/app/persona_read_model.py index 3038b88..65a174a 100644 --- a/apps/api/app/persona_read_model.py +++ b/apps/api/app/persona_read_model.py @@ -144,6 +144,16 @@ def _first_text_value(data: dict[str, Any]) -> str: return "" +def _presenting_summary(data: dict[str, Any]) -> str: + """Return the authored presenting complaint, independent of JSON key order.""" + + for key in ("complaint", "주호소", "presenting_complaint", "chief_complaint"): + value = data.get(key) + if isinstance(value, str) and value.strip(): + return value.strip() + return _first_text_value(data) + + def persona_summary(entry: CatalogPersona) -> PersonaSummary: card = entry.card return PersonaSummary( @@ -155,7 +165,7 @@ def persona_summary(entry: CatalogPersona) -> PersonaSummary: difficulty=card.difficulty, theory_target=card.theory_target, demographics=card.demographics, - presenting_summary=_first_text_value(card.presenting), + presenting_summary=_presenting_summary(card.presenting), source=entry.source, degraded=entry.degraded, ) diff --git a/apps/api/app/routes/admin.py b/apps/api/app/routes/admin.py index ff7c8da..cf2eadd 100644 --- a/apps/api/app/routes/admin.py +++ b/apps/api/app/routes/admin.py @@ -344,6 +344,7 @@ class AdminUserResponse(BaseModel): display_name: str role: RoleName admin_access: bool + learner_feedback_enabled: bool super_admin: bool = False account_status: AccountStatus cohort_ids: list[str] @@ -364,6 +365,7 @@ class AdminUserPatch(BaseModel): display_name: str | None = Field(default=None, min_length=1, max_length=80) role: RoleName | None = None admin_access: bool | None = None + learner_feedback_enabled: bool | None = None account_status: AccountStatus | None = None affiliation: str | None = Field(default=None, max_length=120) cohort_ids: list[str] | None = None @@ -1182,11 +1184,18 @@ def _usage_from_runtime_store(window_days: int) -> AdminUsageResponse: class AdminUserCreate(BaseModel): - email: str = Field(..., min_length=3, max_length=254) + email: str = Field( + ..., + min_length=3, + max_length=254, + pattern=r"^[^@\s]+@[^@\s]+\.[^@\s]+$", + ) display_name: str = Field(..., min_length=1, max_length=80) role: RoleName = "learner" admin_access: bool = False - account_status: AccountStatus = "approved" + learner_feedback_enabled: bool = True + # 외부 연구참여자는 exact-email 사전등록 뒤 별도 승인을 거치게 한다. + account_status: Literal["pending"] = "pending" affiliation: str | None = Field(default=None, max_length=120) cohort_ids: list[str] = Field(default_factory=list) @@ -1574,6 +1583,7 @@ async def _admin_user_response(user, *, durable: bool) -> AdminUserResponse: display_name=user.display_name, role=user.role, admin_access=effective_admin_access, + learner_feedback_enabled=user.learner_feedback_enabled, super_admin=is_super_admin_email(user.email), account_status=user.account_status, cohort_ids=user.cohort_ids, @@ -2171,6 +2181,7 @@ async def create_user( display_name=body.display_name, role=body.role, admin_access=body.admin_access, + learner_feedback_enabled=body.learner_feedback_enabled, account_status=body.account_status, affiliation=body.affiliation, cohort_ids=body.cohort_ids, @@ -2207,6 +2218,7 @@ async def patch_user( display_name=body.display_name, role=body.role, admin_access=body.admin_access, + learner_feedback_enabled=body.learner_feedback_enabled, account_status=body.account_status, affiliation=body.affiliation, cohort_ids=body.cohort_ids, diff --git a/apps/api/app/routes/auth.py b/apps/api/app/routes/auth.py index 52a6f09..395b2f3 100644 --- a/apps/api/app/routes/auth.py +++ b/apps/api/app/routes/auth.py @@ -221,6 +221,25 @@ def validate_google_identity_domain( return normalized_email +def validate_google_identity( + *, + email: str | None, + email_verified: bool, +) -> str: + """Accept every Google account whose email claim is present and verified. + + Google has already validated the account before issuing the ID token. The + application deliberately does not impose an email-domain or pre-registration + gate on top of that provider identity. + """ + normalized_email = _normalize_email(email) + if not normalized_email or not _email_domain(normalized_email): + raise HTTPException(status.HTTP_403_FORBIDDEN, detail="email claim is required") + if not email_verified: + raise HTTPException(status.HTTP_403_FORBIDDEN, detail="email is not verified") + return normalized_email + + async def validate_login_identity_email( *, email: str | None, @@ -774,7 +793,10 @@ async def auth_config(request: Request) -> AuthConfigResponse: google_oauth_configured=google_ready, saml_configured=saml_ready, providers=_auth_provider_statuses(), - allowed_email_domains=sorted(allowed_email_domains()), + # Google OIDC accepts every provider-verified email. Keep the legacy + # setting for dev-login/SAML policy, but do not advertise it as a Google + # restriction to the browser. + allowed_email_domains=[], redirect_uri=settings.oauth_redirect_uri, dev_login_enabled=_dev_login_available(request), ) @@ -937,19 +959,18 @@ async def callback( return _oauth_callback_error("issuer_mismatch", request) try: - email, managed_user = await validate_login_identity_email( + email = validate_google_identity( email=claims.get("email"), email_verified=claims.get("email_verified") in {True, "true", "True", "1", 1}, - hosted_domain=claims.get("hd"), ) except HTTPException: _log_oauth_callback_failure( request, - "domain_not_allowed", + "id_token_invalid", email_domain=_email_domain(str(claims.get("email") or "")), - hosted_domain=_normalize_domain(str(claims.get("hd") or "")), ) - return _oauth_callback_error("domain_not_allowed", request) + return _oauth_callback_error("id_token_invalid", request) + managed_user = await get_managed_user_by_email(email) role = _role_for_managed_user(managed_user, _role_for_email(email)) display_name = str(claims.get("name") or email) cohort_ids = _cohort_ids_for_managed_user( @@ -967,6 +988,7 @@ async def callback( role=role.value, cohort_ids=cohort_ids, external_id=external_id, + account_status="approved", ) except InactiveUserError: _log_oauth_callback_failure( diff --git a/apps/api/app/routes/calibration_transfer.py b/apps/api/app/routes/calibration_transfer.py index 6ee8444..e8024aa 100644 --- a/apps/api/app/routes/calibration_transfer.py +++ b/apps/api/app/routes/calibration_transfer.py @@ -24,7 +24,11 @@ from ..contracts.calibration_transfer import ( ) from ..config import Settings, get_settings from ..deps import AIView, Principal, Role, db_for_ai_view, require_role -from ..services import calibration_transfer_store, session_learning_producer +from ..services import ( + calibration_transfer_store, + feedback_policy, + session_learning_producer, +) router = APIRouter(tags=["calibration-transfer"]) @@ -337,8 +341,8 @@ class ActualTransferExecutionRequest(BaseModel): class ActualTransferExecutionResponse(BaseModel): - execution: ActualTransferExecution - assessment: ActualTransferAssessment + execution: ActualTransferExecution | None = None + assessment: ActualTransferAssessment | None = None idempotent_replay: bool @@ -532,6 +536,22 @@ class CalibrationTransferReadModelResponse(BaseModel): ) +def _learner_input_only_payload(payload: dict[str, Any]) -> dict[str, Any]: + """자기예측 원문만 남기고 내부·교수자·AI 파생 판정을 제거한다.""" + + redacted = dict(payload) + redacted["prediction_histories"] = [ + {**dict(history), "external_observation": None} + for history in payload.get("prediction_histories", []) + ] + redacted["calibration_assessments"] = [] + redacted["transfer_suites"] = [] + redacted["teacher_reviews"] = [] + redacted["actual_executions"] = [] + redacted["actual_transfer_assessments"] = [] + return redacted + + def _http_error(exc: Exception) -> HTTPException: if isinstance( exc, calibration_transfer_store.CalibrationTransferNotFoundError @@ -673,12 +693,20 @@ async def create_actual_transfer_execution( body: ActualTransferExecutionRequest, principal: LearnerPrincipal, ) -> ActualTransferExecutionResponse: + expose_feedback = await feedback_policy.can_expose_session_learner_feedback( + body.practice_session_id, + principal, + ) try: payload = await calibration_transfer_store.append_actual_transfer_execution( principal=principal, **body.model_dump() ) except _STORE_ERRORS as exc: raise _http_error(exc) from exc + if not expose_feedback: + return ActualTransferExecutionResponse( + idempotent_replay=bool(payload.get("idempotent_replay", False)) + ) return ActualTransferExecutionResponse.model_validate(payload) @@ -706,6 +734,7 @@ async def create_teacher_review( ) async def get_my_calibration_transfer( principal: LearnerPrincipal, + session_id: UUID | None = None, ) -> CalibrationTransferReadModelResponse: try: payload = await calibration_transfer_store.read_calibration_transfer( @@ -713,6 +742,23 @@ async def get_my_calibration_transfer( ) except _STORE_ERRORS as exc: raise _http_error(exc) from exc + payload = dict(payload) + snapshot_enabled = bool( + payload.pop("_learner_feedback_snapshot_enabled", True) + ) + expose_feedback = bool( + principal.learner_feedback_enabled and snapshot_enabled + ) + if session_id is not None: + expose_feedback = bool( + expose_feedback + and await feedback_policy.can_expose_session_learner_feedback( + session_id, + principal, + ) + ) + if not expose_feedback: + payload = _learner_input_only_payload(payload) return CalibrationTransferReadModelResponse.model_validate(payload) diff --git a/apps/api/app/routes/deliberate_practices.py b/apps/api/app/routes/deliberate_practices.py index d575749..f4ab2b2 100644 --- a/apps/api/app/routes/deliberate_practices.py +++ b/apps/api/app/routes/deliberate_practices.py @@ -21,7 +21,7 @@ from ..contracts.deliberate_practice import ( PracticePrescription, ) from ..deps import AIView, Principal, Role, db_for_ai_view, require_role -from ..services import deliberate_practice_store +from ..services import deliberate_practice_store, feedback_policy router = APIRouter(tags=["deliberate-practice"]) @@ -248,6 +248,14 @@ class DeliberatePracticeReadModelResponse(BaseModel): def _http_error(exc: Exception) -> HTTPException: + if isinstance( + exc, + deliberate_practice_store.DeliberatePracticeFeedbackDisabledError, + ): + return HTTPException( + status.HTTP_403_FORBIDDEN, + detail=feedback_policy.LEARNER_FEEDBACK_DISABLED_DETAIL, + ) if isinstance(exc, deliberate_practice_store.DeliberatePracticeNotFoundError): return HTTPException(status.HTTP_404_NOT_FOUND, detail=str(exc)) if isinstance(exc, deliberate_practice_store.DeliberatePracticeConflictError): @@ -295,6 +303,7 @@ async def create_practice_attempt( body: PracticeAttemptSubmissionRequest, principal: LearnerPrincipal, ) -> PracticeAttemptSubmissionResponse: + feedback_policy.require_principal_learner_feedback(principal) try: payload = await deliberate_practice_store.append_learner_attempt_submission( principal=principal, @@ -305,6 +314,7 @@ async def create_practice_attempt( except ( deliberate_practice_store.DeliberatePracticeNotFoundError, deliberate_practice_store.DeliberatePracticeConflictError, + deliberate_practice_store.DeliberatePracticeFeedbackDisabledError, deliberate_practice_store.DeliberatePracticeStateError, ) as exc: raise _http_error(exc) from exc @@ -321,6 +331,10 @@ async def observe_completed_practice_session( practice_session_id: UUID, principal: LearnerPrincipal, ) -> PracticeAttemptSubmissionResponse: + await feedback_policy.require_session_learner_feedback( + practice_session_id, + principal, + ) try: payload = await deliberate_practice_store.append_runtime_practice_session( principal=principal, @@ -330,6 +344,7 @@ async def observe_completed_practice_session( except ( deliberate_practice_store.DeliberatePracticeNotFoundError, deliberate_practice_store.DeliberatePracticeConflictError, + deliberate_practice_store.DeliberatePracticeFeedbackDisabledError, deliberate_practice_store.DeliberatePracticeStateError, ) as exc: raise _http_error(exc) from exc @@ -368,6 +383,7 @@ async def correct_practice_attempt( async def get_my_deliberate_practice( principal: LearnerPrincipal, ) -> DeliberatePracticeReadModelResponse: + feedback_policy.require_principal_learner_feedback(principal) try: payload = await deliberate_practice_store.read_deliberate_practice( principal=principal @@ -375,6 +391,7 @@ async def get_my_deliberate_practice( except ( deliberate_practice_store.DeliberatePracticeNotFoundError, deliberate_practice_store.DeliberatePracticeConflictError, + deliberate_practice_store.DeliberatePracticeFeedbackDisabledError, deliberate_practice_store.DeliberatePracticeStateError, ) as exc: raise _http_error(exc) from exc diff --git a/apps/api/app/routes/eval.py b/apps/api/app/routes/eval.py index 40cf2d8..f3da643 100644 --- a/apps/api/app/routes/eval.py +++ b/apps/api/app/routes/eval.py @@ -123,7 +123,13 @@ async def reevaluate_session( 엔진 장애는 503 으로 변환(평가는 비치명적이지만 트리거는 사용자 명시 요청이라 에러 노출). """ sess = await _load_session_or_404(session_id, principal) - enriched = enriched_masked_turns(sess.masked_turns()) + counselor_identity = getattr(sess, "learner_label", None) + client_identity = getattr(sess.persona, "display_name", None) + enriched = enriched_masked_turns( + sess.masked_turns(), + counselor_identity=counselor_identity, + client_identity=client_identity, + ) # 누적 기법 코드 — DB 미가용이라 fast 결과가 없으면 빈 분포(deep LLM 정성 평가는 그대로 유효). technique_codes: list[str] = [] @@ -147,6 +153,8 @@ async def reevaluate_session( scope=body.scope if body.scope in ("session_end", "stage_transition") else "session_end", stage=sess.state.stage.value, error=detail, + counselor_identity=counselor_identity, + client_identity=client_identity, ) saved = await session_persistence.save_session_evaluation(write) if not saved: @@ -161,6 +169,8 @@ async def reevaluate_session( session_id=session_id, learner_id=sess.learner_id, result=result, + counselor_identity=counselor_identity, + client_identity=client_identity, ) saved = await session_persistence.save_session_evaluation(write) if not saved: diff --git a/apps/api/app/routes/kb.py b/apps/api/app/routes/kb.py index 7fe91d2..9e4ce71 100644 --- a/apps/api/app/routes/kb.py +++ b/apps/api/app/routes/kb.py @@ -316,7 +316,8 @@ async def index_document( """문서 인덱싱(관리자, RBAC ADMIN 강제). content_hash 증분 + 청크 임베딩 적재. ⚠️ 임베딩은 무거운 작업 → 본래 BackgroundTasks/배치 워커 위임 권장(202 Accepted). - DSM verbatim 저작권(license C/D)은 source 등록 시점 external_llm_ok 가드 책임. + source_id는 사전 등록된 kb.source만 허용하고, 라이선스와 프로토콜 active 상태는 + 요청값이 아닌 DB 행으로 검증한다. 모델 미가용 시 embedding NULL 폴백(BM25 만, degraded=True) — 크래시 X. """ req = rag.IndexRequest( @@ -329,6 +330,7 @@ async def index_document( try: # 관리자 인덱싱은 RLS 미적용(쓰기 — kb 스키마 직접). role 주입 없이 acquire. async with acquire() as conn: + await rag.validate_index_source(conn, body.source_id) result = await rag.index_document(conn, req) except rag.IndexPolicyViolation as e: raise HTTPException(status.HTTP_422_UNPROCESSABLE_ENTITY, detail=str(e)) from e diff --git a/apps/api/app/routes/measurements.py b/apps/api/app/routes/measurements.py index 0ef0f8d..8118bb5 100644 --- a/apps/api/app/routes/measurements.py +++ b/apps/api/app/routes/measurements.py @@ -24,7 +24,7 @@ from ..contracts.measurement import ( SourceKind, ) from ..deps import CurrentPrincipal, Principal, Role, require_role -from ..services import alliance_measurement +from ..services import alliance_measurement, feedback_policy router = APIRouter(prefix="/sessions", tags=["measurements"]) @@ -168,6 +168,10 @@ async def get_alliance_pulses( session_id: UUID, principal: CurrentPrincipal, ) -> AlliancePulseListResponse: + expose_feedback = await feedback_policy.can_expose_session_learner_feedback( + session_id, + principal, + ) try: items = await alliance_measurement.list_alliance_pulses( principal=principal, @@ -175,6 +179,8 @@ async def get_alliance_pulses( ) except alliance_measurement.AlliancePulseNotFoundError as exc: raise _measurement_http_error(exc) from exc + if not expose_feedback: + items = [{**dict(item), "measurements": []} for item in items] return AlliancePulseListResponse.model_validate({"items": items}) diff --git a/apps/api/app/routes/multimodal_alliance.py b/apps/api/app/routes/multimodal_alliance.py index c176081..b90c1ef 100644 --- a/apps/api/app/routes/multimodal_alliance.py +++ b/apps/api/app/routes/multimodal_alliance.py @@ -23,7 +23,7 @@ from ..contracts.multimodal_alliance import ( ModalityAxisMeasurement, ) from ..deps import AIView, Principal, Role, db_for_ai_view, require_role -from ..services import multimodal_alliance_store +from ..services import feedback_policy, multimodal_alliance_store router = APIRouter(tags=["multimodal-alliance"]) @@ -521,7 +521,14 @@ async def sweep_multimodal_retention( async def get_multimodal_session_metadata( session_id: UUID, principal: HumanPrincipal, + include_derived: bool = True, ) -> MultimodalSessionMetadataResponse: + expose_feedback = await feedback_policy.can_expose_session_learner_feedback( + session_id, + principal, + ) + if include_derived and not expose_feedback: + await feedback_policy.require_session_learner_feedback(session_id, principal) try: payload = await multimodal_alliance_store.read_session_metadata( principal=principal, @@ -529,6 +536,15 @@ async def get_multimodal_session_metadata( ) except multimodal_alliance_store.MultimodalAllianceError as exc: raise _http_error(exc) from exc + if not include_derived: + payload = { + **payload, + "timelines": [], + "word_timestamps": [], + "voice_events": [], + "measurements": [], + "fusion_decisions": [], + } return MultimodalSessionMetadataResponse.model_validate(payload) diff --git a/apps/api/app/routes/outcome_trajectories.py b/apps/api/app/routes/outcome_trajectories.py index 0c0c189..f1ea81f 100644 --- a/apps/api/app/routes/outcome_trajectories.py +++ b/apps/api/app/routes/outcome_trajectories.py @@ -19,7 +19,7 @@ from ..contracts.outcome_trajectory import ( SyntheticExpectedDistribution, ) from ..deps import CurrentPrincipal, Principal, Role, require_role -from ..services import outcome_trajectory_store +from ..services import feedback_policy, outcome_trajectory_store router = APIRouter(prefix="/sessions", tags=["outcome-trajectories"]) @@ -190,6 +190,7 @@ async def get_outcome_trajectory( session_id: UUID, principal: CurrentPrincipal, ) -> OutcomeTrajectoryResponse: + await feedback_policy.require_session_learner_feedback(session_id, principal) try: payload = await outcome_trajectory_store.read_outcome_trajectory( principal=principal, @@ -214,6 +215,7 @@ async def recompute_outcome_trajectory( body: OutcomeTrajectoryRecomputeRequest, principal: CurrentPrincipal, ) -> OutcomeTrajectoryResponse: + await feedback_policy.require_session_learner_feedback(session_id, principal) try: payload = await outcome_trajectory_store.read_outcome_trajectory( principal=principal, diff --git a/apps/api/app/routes/protocols.py b/apps/api/app/routes/protocols.py new file mode 100644 index 0000000..4a056c8 --- /dev/null +++ b/apps/api/app/routes/protocols.py @@ -0,0 +1,220 @@ +"""관리자 전용 상담 프로토콜 등록·활성화·퇴역 API.""" + +from __future__ import annotations + +from datetime import datetime +from typing import Annotated, Literal +from uuid import UUID + +from fastapi import APIRouter, Depends, HTTPException, Query, status +from pydantic import BaseModel, Field, field_validator, model_validator + +from ..db import acquire +from ..deps import Principal, Role, require_role +from ..services import protocol_registry + +router = APIRouter(prefix="/admin/protocols", tags=["admin-protocols"]) +AdminPrincipal = Annotated[Principal, Depends(require_role(Role.ADMIN))] + + +class AdminProtocolCreate(BaseModel): + title: str = Field(..., min_length=1, max_length=240) + source: str = Field(..., min_length=1, max_length=1000) + version: int = Field(default=1, ge=1, le=1_000_000) + license: Literal["A", "B", "C", "D"] + external_llm_ok: bool = False + content: str = Field(..., min_length=1, max_length=500_000) + + @field_validator("title", "source", "content") + @classmethod + def reject_blank_text(cls, value: str) -> str: + if not value.strip(): + raise ValueError("빈 값은 등록할 수 없습니다.") + return value + + @model_validator(mode="after") + def enforce_license_boundary(self) -> "AdminProtocolCreate": + if self.license in {"C", "D"} and self.external_llm_ok: + raise ValueError("라이선스 C/D는 외부 LLM 사용을 허용할 수 없습니다.") + return self + + +class AdminProtocolResponse(BaseModel): + protocol_id: str + source_id: str + title: str + source: str + version: int + license: Literal["A", "B", "C", "D"] + external_llm_ok: bool + content: str + content_hash: str + status: Literal["draft", "active", "retired"] + registered_by: str + registered_at: datetime + activated_at: datetime | None = None + retired_at: datetime | None = None + + +class AdminProtocolListResponse(BaseModel): + protocols: list[AdminProtocolResponse] + total: int + + +class AdminProtocolActivationResponse(BaseModel): + protocol: AdminProtocolResponse + chunks_indexed: int + skipped_unchanged: bool + embedded: bool + degraded: bool + + +def _response(record: protocol_registry.ProtocolRecord) -> AdminProtocolResponse: + return AdminProtocolResponse( + protocol_id=record.protocol_id, + source_id=record.source_id, + title=record.title, + source=record.source, + version=record.version, + license=record.license, + external_llm_ok=record.external_llm_ok, + content=record.content, + content_hash=record.content_hash, + status=record.status, + registered_by=record.registered_by, + registered_at=record.registered_at, + activated_at=record.activated_at, + retired_at=record.retired_at, + ) + + +def _raise_http(error: protocol_registry.ProtocolRegistryError) -> None: + if isinstance(error, protocol_registry.ProtocolNotFound): + raise HTTPException(status.HTTP_404_NOT_FOUND, detail=str(error)) from error + if isinstance(error, protocol_registry.ProtocolTransitionConflict): + raise HTTPException(status.HTTP_409_CONFLICT, detail=str(error)) from error + if isinstance(error, protocol_registry.ProtocolPolicyViolation): + raise HTTPException(status.HTTP_422_UNPROCESSABLE_ENTITY, detail=str(error)) from error + raise HTTPException(status.HTTP_503_SERVICE_UNAVAILABLE, detail=str(error)) from error + + +@router.get("", response_model=AdminProtocolListResponse) +async def list_admin_protocols( + principal: AdminPrincipal, + status_filter: Annotated[ + Literal["draft", "active", "retired"] | None, + Query(alias="status"), + ] = None, + search: Annotated[str | None, Query(max_length=240)] = None, +) -> AdminProtocolListResponse: + try: + async with acquire(role="admin", user_id=principal.user_id) as conn: + records = await protocol_registry.list_protocols( + conn, + status_filter=status_filter, + search=search, + ) + except protocol_registry.ProtocolRegistryError as error: + _raise_http(error) + except RuntimeError as error: + raise HTTPException( + status.HTTP_503_SERVICE_UNAVAILABLE, + detail=f"프로토콜 저장소를 사용할 수 없습니다: {error}", + ) from error + return AdminProtocolListResponse( + protocols=[_response(record) for record in records], + total=len(records), + ) + + +@router.post( + "", + response_model=AdminProtocolResponse, + status_code=status.HTTP_201_CREATED, +) +async def create_admin_protocol( + body: AdminProtocolCreate, + principal: AdminPrincipal, +) -> AdminProtocolResponse: + try: + async with acquire(role="admin", user_id=principal.user_id) as conn: + record = await protocol_registry.create_protocol( + conn, + title=body.title, + source=body.source, + version=body.version, + license_class=body.license, + external_llm_ok=body.external_llm_ok, + content=body.content, + registered_by=principal.user_id, + ) + except protocol_registry.ProtocolRegistryError as error: + _raise_http(error) + except RuntimeError as error: + raise HTTPException( + status.HTTP_503_SERVICE_UNAVAILABLE, + detail=f"프로토콜 저장소를 사용할 수 없습니다: {error}", + ) from error + return _response(record) + + +@router.post( + "/{protocol_id}/activate", + response_model=AdminProtocolActivationResponse, +) +async def activate_admin_protocol( + protocol_id: UUID, + principal: AdminPrincipal, +) -> AdminProtocolActivationResponse: + try: + async with acquire(role="admin", user_id=principal.user_id) as conn: + async with conn.transaction(): + record, indexed = await protocol_registry.activate_protocol( + conn, + protocol_id=str(protocol_id), + ) + except protocol_registry.ProtocolRegistryError as error: + _raise_http(error) + except RuntimeError as error: + raise HTTPException( + status.HTTP_503_SERVICE_UNAVAILABLE, + detail=f"프로토콜 저장소를 사용할 수 없습니다: {error}", + ) from error + return AdminProtocolActivationResponse( + protocol=_response(record), + chunks_indexed=indexed.chunks_indexed, + skipped_unchanged=indexed.skipped_unchanged, + embedded=indexed.embedded, + degraded=indexed.degraded, + ) + + +@router.post("/{protocol_id}/retire", response_model=AdminProtocolResponse) +async def retire_admin_protocol( + protocol_id: UUID, + principal: AdminPrincipal, +) -> AdminProtocolResponse: + try: + async with acquire(role="admin", user_id=principal.user_id) as conn: + async with conn.transaction(): + record = await protocol_registry.retire_protocol( + conn, + protocol_id=str(protocol_id), + ) + except protocol_registry.ProtocolRegistryError as error: + _raise_http(error) + except RuntimeError as error: + raise HTTPException( + status.HTTP_503_SERVICE_UNAVAILABLE, + detail=f"프로토콜 저장소를 사용할 수 없습니다: {error}", + ) from error + return _response(record) + + +__all__ = [ + "AdminProtocolActivationResponse", + "AdminProtocolCreate", + "AdminProtocolListResponse", + "AdminProtocolResponse", + "router", +] diff --git a/apps/api/app/routes/rupture_repairs.py b/apps/api/app/routes/rupture_repairs.py index 6e02c8b..a49d199 100644 --- a/apps/api/app/routes/rupture_repairs.py +++ b/apps/api/app/routes/rupture_repairs.py @@ -15,7 +15,7 @@ from pydantic import BaseModel, ConfigDict, Field, field_validator, model_valida from ..contracts.rupture_repair import RuptureLifecycleState, RuptureType from ..config import Settings, get_settings from ..deps import AIView, CurrentPrincipal, db_for_ai_view -from ..services import rupture_repair_store +from ..services import feedback_policy, rupture_repair_store router = APIRouter(tags=["rupture-repairs"]) @@ -358,6 +358,7 @@ async def get_rupture_repairs( session_id: UUID, principal: CurrentPrincipal, ) -> RuptureRepairReadModelResponse: + await feedback_policy.require_session_learner_feedback(session_id, principal) try: payload = await rupture_repair_store.read_rupture_repairs( principal=principal, @@ -435,6 +436,7 @@ async def create_human_rupture_correction( body: HumanRuptureCorrectionRequest, principal: CurrentPrincipal, ) -> HumanRuptureCorrectionResponse: + await feedback_policy.require_session_learner_feedback(session_id, principal) try: observation_id = await rupture_repair_store.append_human_correction( principal=principal, diff --git a/apps/api/app/routes/sessions.py b/apps/api/app/routes/sessions.py index 3e23100..f6bce87 100644 --- a/apps/api/app/routes/sessions.py +++ b/apps/api/app/routes/sessions.py @@ -30,6 +30,7 @@ from ..runtime_policy import require_runtime_fallback_allowed from ..session_evaluation_input import enriched_masked_turns from ..services import ( evaluator, + feedback_policy, guardrail, live_coach, memory, @@ -67,6 +68,7 @@ from ..session_read_model import ( dashboard_growth as _dashboard_growth, dashboard_overview as _dashboard_overview, dashboard_persona_progress as _dashboard_persona_progress, + dashboard_training_exposure as _dashboard_training_exposure, iso as _iso, learner_summary as _learner_summary, learner_visible_turns as _learner_visible_turns, @@ -112,6 +114,8 @@ class SessionStartRequest(BaseModel): class SessionStartResponse(BaseModel): session_id: str case_id: str + persona_id: str + persona_version: int session_no: int stage: StageLabel effective_openness: float @@ -122,6 +126,7 @@ class SessionStartResponse(BaseModel): # 시간 기반 회기 종료 계약(회의 P1): 프론트 타이머·10분 전 알람의 기준값. duration_limit_seconds: int = 0 warning_before_end_seconds: int = 0 + learner_feedback_enabled: bool = True class TurnRequest(BaseModel): @@ -308,6 +313,8 @@ async def _retrieve_live_coach_grounding( source_type=source_type or None, version=source_version or None, citation=citation or None, + license_class=chunk.license_class, + external_llm_ok=chunk.external_llm_ok, summary=body[:500], ) ) @@ -525,6 +532,7 @@ async def _prepare_turn_context( card=sess.persona, state=sess.state, learner_text=learner_text, + learner_identity=sess.learner_label, memory=orchestrator.TurnMemory( recall_summary=recall.recall_summary, pinned_facts=recall.pinned_facts, @@ -908,7 +916,11 @@ async def _generate_and_save_session_evaluation(sess: InProcSession) -> None: return timeout_seconds = _session_evaluation_timeout_seconds() - enriched = enriched_masked_turns(sess.masked_turns()) + enriched = enriched_masked_turns( + sess.masked_turns(), + counselor_identity=sess.learner_label, + client_identity=sess.persona.display_name, + ) try: result = await asyncio.wait_for( @@ -928,6 +940,8 @@ async def _generate_and_save_session_evaluation(sess: InProcSession) -> None: session_id=sess.session_id, learner_id=sess.learner_id, result=result, + counselor_identity=sess.learner_label, + client_identity=sess.persona.display_name, ) saved = await session_persistence.save_session_evaluation(write) if not saved: @@ -960,6 +974,8 @@ async def _generate_and_save_session_evaluation(sess: InProcSession) -> None: scope="session_end", stage=_stage_label(sess.state.stage), error=message, + counselor_identity=sess.learner_label, + client_identity=sess.persona.display_name, ) saved = await session_persistence.save_session_evaluation(write) if not saved: @@ -978,6 +994,8 @@ async def _generate_and_save_session_evaluation(sess: InProcSession) -> None: scope="session_end", stage=_stage_label(sess.state.stage), error=exc, + counselor_identity=sess.learner_label, + client_identity=sess.persona.display_name, ) saved = await session_persistence.save_session_evaluation(write) if not saved: @@ -1124,6 +1142,8 @@ async def _load_learner_sessions( async def _review_ready(sess: InProcSession, principal: Principal) -> bool: + if not feedback_policy.can_expose_principal_learner_feedback(sess, principal): + return False turns = _learner_visible_turns(sess) if not sess.ended or not turns: return False @@ -1173,6 +1193,9 @@ async def _session_archive_response( session=_learner_summary( sess, review_ready=await _review_ready(sess, principal), + learner_feedback_enabled=( + feedback_policy.effective_learner_feedback_enabled(sess, principal) + ), archived=archived, archived_at=archived_at, ), @@ -1194,6 +1217,12 @@ async def list_learner_sessions(principal: CurrentPrincipal) -> LearnerSessionsR _learner_summary( sess, review_ready=await _review_ready(sess, principal), + learner_feedback_enabled=( + feedback_policy.effective_learner_feedback_enabled( + sess, + principal, + ) + ), archived=archive_record is not None, archived_at=_iso(float(archived_at)) if isinstance(archived_at, (int, float)) @@ -1213,10 +1242,22 @@ async def learner_dashboard(principal: CurrentPrincipal) -> LearnerDashboardResp principal = _ensure_learner(principal) sessions, durable = await _load_learner_sessions( principal, - include_turn_evaluation=True, + include_turn_evaluation=principal.learner_feedback_enabled, ) review_ready = await _review_ready_map(sessions, principal) archives = await _archive_map(sessions, principal) + learner_feedback_enabled = { + sess.session_id: feedback_policy.effective_learner_feedback_enabled( + sess, + principal, + ) + for sess in sessions + } + feedback_sessions = [ + sess + for sess in sessions + if learner_feedback_enabled.get(sess.session_id, True) + ] visible_review_ready = { session_id: ready for session_id, ready in review_ready.items() @@ -1229,10 +1270,15 @@ async def learner_dashboard(principal: CurrentPrincipal) -> LearnerDashboardResp visible_review_ready=visible_review_ready, archived_sessions=len(archives), ), - growth=_dashboard_growth(sessions), - persona_progress=_dashboard_persona_progress(sessions, visible_review_ready), + growth=_dashboard_growth(feedback_sessions), + persona_progress=_dashboard_persona_progress( + sessions, + visible_review_ready, + learner_feedback_enabled, + ), + training_exposure=_dashboard_training_exposure(sessions), achievements=_dashboard_achievements(sessions, visible_review_ready), - recent_feedback=_dashboard_feedback(sessions), + recent_feedback=_dashboard_feedback(feedback_sessions), message=( "실제 연습 기록을 기준으로 개인 학습 흐름을 표시합니다." if sessions @@ -1253,7 +1299,13 @@ async def get_session_detail( principal, allow_ended=True, ) - return _session_detail(sess, review_ready=await _review_ready(sess, principal)) + return _session_detail( + sess, + review_ready=await _review_ready(sess, principal), + learner_feedback_enabled=( + feedback_policy.effective_learner_feedback_enabled(sess, principal) + ), + ) @router.post("/{session_id}/archive", response_model=SessionArchiveResponse) @@ -1355,18 +1407,26 @@ async def start_session( carry_rapport = st.rapport_credit goal_stages = [str(stage) for stage in body.goal_stages] - sess = await session_persistence.create_session( - learner_id=principal.user_id, - card=card, - theory_mode=body.theory_mode, - state=st, - session_no=session_no, - carry_rapport=carry_rapport, - persona_id=catalog_persona.persona_id, - persona_version=catalog_persona.version, - case_id=case_context.case_id if case_context else None, - goal_stages=goal_stages, - ) + learner_feedback_enabled = principal.learner_feedback_enabled + try: + sess = await session_persistence.create_session( + learner_id=principal.user_id, + card=card, + theory_mode=body.theory_mode, + state=st, + session_no=session_no, + carry_rapport=carry_rapport, + persona_id=catalog_persona.persona_id, + persona_version=catalog_persona.version, + case_id=case_context.case_id if case_context else None, + goal_stages=goal_stages, + learner_feedback_enabled=learner_feedback_enabled, + ) + except session_persistence.SessionCreationPersistenceError as exc: + raise HTTPException( + status.HTTP_503_SERVICE_UNAVAILABLE, + detail="session_persistence_unavailable", + ) from exc degraded = catalog_persona.degraded or sess is None if sess is None: require_runtime_fallback_allowed("session creation") @@ -1375,13 +1435,21 @@ async def start_session( persona=card, theory_mode=body.theory_mode, state=st, + persona_id=catalog_persona.persona_id, + persona_version=catalog_persona.version, session_no=session_no, carry_rapport=carry_rapport, goal_stages=goal_stages, + learner_feedback_enabled=learner_feedback_enabled, ) else: store.put(sess) + # DB 재조회 전의 첫 턴과 runtime fallback에서도 인증된 학습자 표시명을 + # 역할 기반 비식별화에 사용할 수 있게 in-process 세션에만 보존한다. + sess.learner_label = principal.display_name or principal.user_id + store.put(sess) + # 즉시 빈/carry 회상으로 응답을 막지 않는다. RAG 회상·KB 단서(임베더 로드 수 초)는 # 백그라운드 warm으로 캐시 — 회기 시작/턴 응답이 임베더 로드에 블로킹되지 않게(성능 회귀 방지). _RECALL_CACHE[sess.session_id] = recall @@ -1390,6 +1458,8 @@ async def start_session( return SessionStartResponse( session_id=sess.session_id, case_id=sess.case_id, + persona_id=catalog_persona.persona_id, + persona_version=catalog_persona.version, session_no=sess.session_no, stage=_stage_label(st.stage), effective_openness=round(st.effective_openness, 4), @@ -1399,6 +1469,7 @@ async def start_session( goal_stages=body.goal_stages, duration_limit_seconds=settings.session_duration_minutes * 60, warning_before_end_seconds=settings.session_warning_minutes * 60, + learner_feedback_enabled=sess.learner_feedback_enabled, ) @@ -1413,18 +1484,31 @@ async def get_session_review( principal, include_turn_evaluation=True, ) - ( - evaluation_record, - evaluation_durable, - ) = await session_persistence.load_session_evaluation( - session_id, - review_principal, + expose_learner_feedback = ( + feedback_policy.can_expose_principal_learner_feedback( + sess, + review_principal, + ) ) + if expose_learner_feedback: + ( + evaluation_record, + evaluation_durable, + ) = await session_persistence.load_session_evaluation( + session_id, + review_principal, + ) + else: + evaluation_record, evaluation_durable = None, True saved_worksheet_payload, _ = await session_persistence.load_case_worksheet( session_id, review_principal, ) include_teacher_review = review_principal.role in {Role.TEACHER, Role.ADMIN} + learner_feedback_enabled = feedback_policy.effective_learner_feedback_enabled( + sess, + review_principal, + ) teacher_review_record = None if include_teacher_review: teacher_review_record, _ = await session_persistence.load_session_review_status( @@ -1440,6 +1524,8 @@ async def get_session_review( saved_worksheet_payload=saved_worksheet_payload, include_teacher_review=include_teacher_review, teacher_review_record=teacher_review_record, + learner_feedback_enabled=learner_feedback_enabled, + expose_learner_feedback=expose_learner_feedback, ) ) @@ -1462,6 +1548,11 @@ async def create_session_share( raise HTTPException( status.HTTP_409_CONFLICT, detail="session must be ended before sharing" ) + if not feedback_policy.can_expose_principal_learner_feedback(sess, principal): + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail=feedback_policy.LEARNER_FEEDBACK_DISABLED_DETAIL, + ) review = await get_session_review(session_id, principal) token = secrets.token_urlsafe(32) @@ -1606,7 +1697,12 @@ async def list_live_coach_history( ) -> LiveCoachHistoryResponse: """현재 회기에서 학습자에게 실제로 전달된 라이브 코칭 이력을 반환한다.""" principal = _ensure_learner(principal) - await _load_session_or_404(session_id, principal) + sess = await _load_session_or_404(session_id, principal) + if not feedback_policy.can_expose_principal_learner_feedback(sess, principal): + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail=feedback_policy.LEARNER_FEEDBACK_DISABLED_DETAIL, + ) events, durable = await session_persistence.list_live_coach_events( session_id, principal ) @@ -1641,6 +1737,11 @@ async def live_coach_turn( """방금 완료된 턴에 대한 비차단 라이브 코칭을 반환한다.""" principal = _ensure_learner(principal) sess = await _load_session_or_404(session_id, principal) + if not feedback_policy.can_expose_principal_learner_feedback(sess, principal): + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail=feedback_policy.LEARNER_FEEDBACK_DISABLED_DETAIL, + ) quota, _ = await session_persistence.get_live_coach_quota(session_id, principal) if int(quota.get("remaining", 0)) <= 0: raise HTTPException( diff --git a/apps/api/app/routes/voice.py b/apps/api/app/routes/voice.py index 5a2f172..84059b5 100644 --- a/apps/api/app/routes/voice.py +++ b/apps/api/app/routes/voice.py @@ -1458,6 +1458,7 @@ async def _prepare_voice_turn_context( card=sess.persona, state=sess.state, learner_text=learner_text, + learner_identity=sess.learner_label, memory=orchestrator.TurnMemory( recall_summary=recall.recall_summary, pinned_facts=recall.pinned_facts, diff --git a/apps/api/app/services/calibration_transfer_store.py b/apps/api/app/services/calibration_transfer_store.py index c6bd4ba..48e8460 100644 --- a/apps/api/app/services/calibration_transfer_store.py +++ b/apps/api/app/services/calibration_transfer_store.py @@ -71,6 +71,14 @@ def _value(row: Mapping[str, Any], key: str, default: Any = None) -> Any: return default +def _public_row(row: Mapping[str, Any]) -> dict[str, Any]: + """API 응답에서 학습자 피드백 정책 판정 전용 열을 제거한다.""" + + payload = dict(row) + payload.pop("source_learner_feedback_enabled", None) + return payload + + def _canonical_hash(payload: Mapping[str, Any]) -> str: serialized = json.dumps( payload, @@ -1544,10 +1552,17 @@ async def read_calibration_transfer( histories = list( await conn.fetch( """ - SELECT history_id, session_id, competency_id, practice_block_id, - scenario_variant_id, phrase_family_id, created_at - FROM app.calibration_prediction_history - WHERE learner_id = $1 ORDER BY created_at, history_id + SELECT history.history_id, history.session_id, + history.competency_id, history.practice_block_id, + history.scenario_variant_id, history.phrase_family_id, + history.created_at, + source_session.learner_feedback_enabled + AS source_learner_feedback_enabled + FROM app.calibration_prediction_history history + JOIN app.sessions source_session + ON source_session.id = history.session_id + WHERE history.learner_id = $1 + ORDER BY history.created_at, history.history_id """, target_learner_id, ) @@ -1615,10 +1630,13 @@ async def read_calibration_transfer( a.source_observation_ids, a.assessment_payload, a.model_run_id, a.instrument_id, a.instrument_version, a.evidence_turn_ids, a.created_at, - p.prescription_id, p.prescription_payload + p.prescription_id, p.prescription_payload, + source_session.learner_feedback_enabled + AS source_learner_feedback_enabled FROM app.calibration_assessment_snapshot a JOIN app.calibration_metacognitive_prescription p ON p.assessment_snapshot_id = a.assessment_snapshot_id + JOIN app.sessions source_session ON source_session.id = a.session_id WHERE a.learner_id = $1 ORDER BY a.competency_id, a.snapshot_no """, @@ -1628,12 +1646,18 @@ async def read_calibration_transfer( suites = list( await conn.fetch( """ - SELECT transfer_suite_record_id, submission_id, suite_key, - session_id, training_phrase_family_ids, model_run_id, - instrument_id, instrument_version, data_classification, - clinical_claim_allowed, created_at - FROM app.calibration_transfer_suite - WHERE learner_id = $1 ORDER BY created_at, transfer_suite_record_id + SELECT suite.transfer_suite_record_id, suite.submission_id, + suite.suite_key, suite.session_id, + suite.training_phrase_family_ids, suite.model_run_id, + suite.instrument_id, suite.instrument_version, + suite.data_classification, suite.clinical_claim_allowed, + suite.created_at, + source_session.learner_feedback_enabled + AS source_learner_feedback_enabled + FROM app.calibration_transfer_suite suite + JOIN app.sessions source_session ON source_session.id = suite.session_id + WHERE suite.learner_id = $1 + ORDER BY suite.created_at, suite.transfer_suite_record_id """, target_learner_id, ) @@ -1715,14 +1739,29 @@ async def read_calibration_transfer( actual_event_rows = list( await conn.fetch( """ - SELECT * - FROM app.calibration_transfer_execution_event - WHERE learner_id = $1 - ORDER BY competency_id, created_at, execution_event_id + SELECT execution.*, + practice_session.learner_feedback_enabled + AS source_learner_feedback_enabled + FROM app.calibration_transfer_execution_event execution + JOIN app.sessions practice_session + ON practice_session.id = execution.practice_session_id + WHERE execution.learner_id = $1 + ORDER BY execution.competency_id, execution.created_at, + execution.execution_event_id """, target_learner_id, ) ) + learner_feedback_snapshot_enabled = all( + bool( + _value( + item, + "source_learner_feedback_enabled", + True, + ) + ) + for item in (*histories, *assessments, *suites, *actual_event_rows) + ) revisions_by_history: dict[UUID, list[dict[str, Any]]] = {} for row in revisions: @@ -1737,7 +1776,7 @@ async def read_calibration_transfer( } prediction_histories: list[dict[str, Any]] = [] for row in histories: - item = dict(row) + item = _public_row(row) history_id = UUID(str(_value(row, "history_id"))) item["revisions"] = revisions_by_history.get(history_id, []) item["lock"] = locks_by_history.get(history_id) @@ -1760,7 +1799,7 @@ async def read_calibration_transfer( ).append(dict(row)) suite_payloads: list[dict[str, Any]] = [] for row in suites: - item = dict(row) + item = _public_row(row) suite_id = UUID(str(_value(row, "transfer_suite_record_id"))) item["trials"] = trials_by_suite.get(suite_id, []) item["assessments"] = assessments_by_suite.get(suite_id, []) @@ -1774,8 +1813,9 @@ async def read_calibration_transfer( "learner_id": target_learner_id, "requested_view": requested_view, "clinical_claim_allowed": False, + "_learner_feedback_snapshot_enabled": learner_feedback_snapshot_enabled, "prediction_histories": prediction_histories, - "calibration_assessments": [dict(item) for item in assessments], + "calibration_assessments": [_public_row(item) for item in assessments], "transfer_suites": suite_payloads, "teacher_reviews": [dict(item) for item in reviews], "actual_executions": [ diff --git a/apps/api/app/services/deliberate_practice_store.py b/apps/api/app/services/deliberate_practice_store.py index 85ded59..3493ac1 100644 --- a/apps/api/app/services/deliberate_practice_store.py +++ b/apps/api/app/services/deliberate_practice_store.py @@ -56,6 +56,10 @@ class DeliberatePracticeConflictError(RuntimeError): pass +class DeliberatePracticeFeedbackDisabledError(PermissionError): + """처방/연습 원천 회기의 학습자 피드백 스냅샷이 비활성이다.""" + + def _value(row: Mapping[str, Any], key: str, default: Any = None) -> Any: try: return row[key] @@ -63,6 +67,14 @@ def _value(row: Mapping[str, Any], key: str, default: Any = None) -> Any: return default +def _public_row(row: Mapping[str, Any]) -> dict[str, Any]: + """API 응답에서 정책 판정 전용 내부 열을 제거한다.""" + + payload = dict(row) + payload.pop("source_learner_feedback_enabled", None) + return payload + + def _canonical_hash(payload: Mapping[str, Any]) -> str: serialized = json.dumps( payload, @@ -202,11 +214,15 @@ async def _latest_snapshot( ) -> Mapping[str, Any] | None: return await conn.fetchrow( """ - SELECT snapshot_id, session_id, snapshot_no, content_hash, graph_payload, - evidence_turn_ids, created_at - FROM app.competency_graph_snapshot - WHERE learner_id = $1 - ORDER BY snapshot_no DESC + SELECT snapshot.snapshot_id, snapshot.session_id, snapshot.snapshot_no, + snapshot.content_hash, snapshot.graph_payload, + snapshot.evidence_turn_ids, snapshot.created_at, + source_session.learner_feedback_enabled + AS source_learner_feedback_enabled + FROM app.competency_graph_snapshot snapshot + JOIN app.sessions source_session ON source_session.id = snapshot.session_id + WHERE snapshot.learner_id = $1 + ORDER BY snapshot.snapshot_no DESC LIMIT 1 """, learner_id, @@ -595,6 +611,37 @@ async def append_learner_attempt_submission( "SELECT pg_advisory_xact_lock(hashtextextended($1::text, 0))", f"practice-attempt:{submission_id}", ) + prescription_row = await conn.fetchrow( + """ + SELECT prescription.prescription_record_id, + prescription.session_id, + prescription.prescription_payload, + prescription.created_at, + source_session.learner_feedback_enabled + AS source_learner_feedback_enabled + FROM app.practice_prescription prescription + JOIN app.sessions source_session + ON source_session.id = prescription.session_id + WHERE prescription.learner_id = $1 + AND prescription.prescription_key = $2 + """, + learner_id, + prescription_id, + ) + if prescription_row is None: + raise DeliberatePracticeNotFoundError( + "practice prescription not found or not visible" + ) + if not bool( + _value( + prescription_row, + "source_learner_feedback_enabled", + True, + ) + ): + raise DeliberatePracticeFeedbackDisabledError( + "learner feedback was disabled for the prescription source session" + ) existing = await _existing_submission( conn, table="app.practice_episode_submission", @@ -608,19 +655,6 @@ async def append_learner_attempt_submission( "SELECT pg_advisory_xact_lock(hashtextextended($1::text, 0))", f"practice-graph:{learner_id}", ) - prescription_row = await conn.fetchrow( - """ - SELECT prescription_record_id, session_id, prescription_payload, created_at - FROM app.practice_prescription - WHERE learner_id = $1 AND prescription_key = $2 - """, - learner_id, - prescription_id, - ) - if prescription_row is None: - raise DeliberatePracticeNotFoundError( - "practice prescription not found or not visible" - ) source_session_id = UUID(str(_value(prescription_row, "session_id"))) session_id = practice_session_id or source_session_id practice_session = await _visible_session(conn, session_id) @@ -1218,10 +1252,13 @@ async def read_deliberate_practice( p.activity_mode, p.scenario_variant_id, p.scenario_novelty, p.difficulty_level, p.prescription_payload, p.created_at, c.card_key, c.coach_claim, c.evidence_turn_ids, c.source_refs, - c.uncertainty, c.counterevidence + c.uncertainty, c.counterevidence, + source_session.learner_feedback_enabled + AS source_learner_feedback_enabled FROM app.practice_prescription p JOIN app.practice_coaching_card c ON c.coaching_card_record_id = p.coaching_card_record_id + JOIN app.sessions source_session ON source_session.id = p.session_id WHERE p.learner_id = $1 ORDER BY p.created_at, p.prescription_record_id """, @@ -1231,12 +1268,19 @@ async def read_deliberate_practice( episodes = list( await conn.fetch( """ - SELECT episode_submission_id, episode_key, session_id, - progress, mastery_allowed, mastery_blockers, uncertainty, - evidence_turn_ids, counterevidence, assessment_payload, created_at - FROM app.practice_episode_submission - WHERE learner_id = $1 - ORDER BY created_at, episode_submission_id + SELECT episode.episode_submission_id, episode.episode_key, + episode.session_id, episode.progress, + episode.mastery_allowed, episode.mastery_blockers, + episode.uncertainty, episode.evidence_turn_ids, + episode.counterevidence, episode.assessment_payload, + episode.created_at, + source_session.learner_feedback_enabled + AS source_learner_feedback_enabled + FROM app.practice_episode_submission episode + JOIN app.sessions source_session + ON source_session.id = episode.session_id + WHERE episode.learner_id = $1 + ORDER BY episode.created_at, episode.episode_submission_id """, target_learner_id, ) @@ -1285,6 +1329,19 @@ async def read_deliberate_practice( else [] ) snapshot = await _latest_snapshot(conn, target_learner_id) + if principal.role == Role.LEARNER and any( + not bool( + _value( + item, + "source_learner_feedback_enabled", + True, + ) + ) + for item in (*prescriptions, *episodes, *([snapshot] if snapshot else [])) + ): + raise DeliberatePracticeFeedbackDisabledError( + "learner feedback was disabled for a practice source session" + ) decision = ( await conn.fetchrow( """ @@ -1313,7 +1370,7 @@ async def read_deliberate_practice( ).append(payload) episode_payloads: list[dict[str, Any]] = [] for row in episodes: - payload = dict(row) + payload = _public_row(row) payload["attempts"] = attempts_by_episode.get( UUID(str(_value(row, "episode_submission_id"))), [] ) @@ -1321,7 +1378,7 @@ async def read_deliberate_practice( return { "learner_id": target_learner_id, "clinical_claim_allowed": False, - "prescriptions": [dict(item) for item in prescriptions], + "prescriptions": [_public_row(item) for item in prescriptions], "episodes": episode_payloads, "competency_graph": ( _value(snapshot, "graph_payload") if snapshot is not None else None @@ -1335,6 +1392,7 @@ async def read_deliberate_practice( __all__ = [ "DeliberatePracticeConflictError", + "DeliberatePracticeFeedbackDisabledError", "DeliberatePracticeNotFoundError", "DeliberatePracticeStateError", "append_learner_attempt_submission", diff --git a/apps/api/app/services/evaluator.py b/apps/api/app/services/evaluator.py index e4949a1..d9f6cb2 100644 --- a/apps/api/app/services/evaluator.py +++ b/apps/api/app/services/evaluator.py @@ -520,8 +520,11 @@ def build_fast_messages(ctx: "TurnContext", client_reply: str) -> list[EngineMes """fast-loop 평가 프롬프트(L0 역할 + 후보 라벨 + 이번 턴 맥락).""" st = ctx.state_after or ctx.state_before theory = _theory_mode(ctx) - client_reply_masked = guardrail.mask_synthetic_generated_pii( - client_reply + client_reply_masked = guardrail.mask_role_identities( + client_reply, + counselor_identity=ctx.counselor_identity, + client_identity=ctx.client_identity, + synthetic_generated=True, ).text_masked recent = ( "\n".join( diff --git a/apps/api/app/services/feedback_policy.py b/apps/api/app/services/feedback_policy.py new file mode 100644 index 0000000..8c1a91c --- /dev/null +++ b/apps/api/app/services/feedback_policy.py @@ -0,0 +1,141 @@ +"""계정별 학습자 AI 피드백 노출 정책의 단일 경계.""" + +from __future__ import annotations + +from typing import Protocol + +from fastapi import HTTPException, status + + +LEARNER_FEEDBACK_DISABLED_DETAIL = "learner_feedback_disabled" + + +class FeedbackPolicySession(Protocol): + learner_feedback_enabled: bool + + +class FeedbackPolicyPrincipal(Protocol): + role: object + learner_feedback_enabled: bool + + +def _role_value(role: object) -> str: + return str(getattr(role, "value", role)) + + +def learner_feedback_enabled(session: FeedbackPolicySession) -> bool: + """구 회기 객체는 호환성을 위해 피드백 허용으로 취급한다.""" + + return bool(getattr(session, "learner_feedback_enabled", True)) + + +def effective_learner_feedback_enabled( + session: FeedbackPolicySession, + principal: FeedbackPolicyPrincipal, +) -> bool: + """학습자는 현재 계정과 회기 스냅샷이 모두 켜져야 피드백이 활성이다.""" + + snapshot_enabled = learner_feedback_enabled(session) + if _role_value(principal.role) != "learner": + return snapshot_enabled + return snapshot_enabled and bool( + getattr(principal, "learner_feedback_enabled", True) + ) + + +def can_expose_principal_learner_feedback( + session: FeedbackPolicySession, + principal: FeedbackPolicyPrincipal, +) -> bool: + """교수자·관리자는 유지하고 학습자에게 effective AND 정책을 적용한다.""" + + return _role_value(principal.role) != "learner" or ( + effective_learner_feedback_enabled(session, principal) + ) + + +def can_expose_learner_feedback( + session: FeedbackPolicySession, + *, + viewer_role: str, +) -> bool: + """관리자·교수자 검토는 유지하고 학습자에게만 스냅샷 정책을 적용한다.""" + + return _role_value(viewer_role) != "learner" or learner_feedback_enabled(session) + + +def require_learner_feedback( + session: FeedbackPolicySession, + *, + viewer_role: object, +) -> None: + """학습자에게만 회기 스냅샷 정책을 적용하고 파생 출력을 fail-closed 한다.""" + + if can_expose_learner_feedback( + session, + viewer_role=_role_value(viewer_role), + ): + return + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail=LEARNER_FEEDBACK_DISABLED_DETAIL, + ) + + +def require_principal_learner_feedback(principal: FeedbackPolicyPrincipal) -> None: + """회기 ID가 없는 학습자 전용 파생 원장은 계정 정책으로 차단한다.""" + + if _role_value(principal.role) != "learner" or bool( + getattr(principal, "learner_feedback_enabled", True) + ): + return + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail=LEARNER_FEEDBACK_DISABLED_DETAIL, + ) + + +async def can_expose_session_learner_feedback( + session_id: object, + principal: FeedbackPolicyPrincipal, +) -> bool: + """현재 계정과 생성 당시 회기 스냅샷을 모두 확인한다. + + 교수자·관리자 감독 화면은 그대로 유지한다. 학습자는 관리자가 현재 피드백을 + 끈 경우 즉시 차단하고, 다시 켠 뒤에도 비활성 상태로 생성된 회기에서는 파생 + 출력을 되살리지 않는다. + """ + + if _role_value(principal.role) != "learner": + return True + if not bool(getattr(principal, "learner_feedback_enabled", True)): + return False + + # 순환 import를 피하면서 세션 조회 경계와 동일한 RLS/런타임 폴백을 사용한다. + from .. import session_persistence + from ..store import store + + session_key = str(session_id) + session = await session_persistence.load_session( + session_key, + principal, # type: ignore[arg-type] + allow_ended=True, + ) + if session is None: + session = store.get(session_key) + if session is None: + # 소유권·존재 오류는 각 도메인 저장소가 원래 계약대로 처리한다. + return True + return can_expose_principal_learner_feedback(session, principal) + + +async def require_session_learner_feedback( + session_id: object, + principal: FeedbackPolicyPrincipal, +) -> None: + if await can_expose_session_learner_feedback(session_id, principal): + return + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail=LEARNER_FEEDBACK_DISABLED_DETAIL, + ) diff --git a/apps/api/app/services/guardrail.py b/apps/api/app/services/guardrail.py index c1d1e31..94bcef0 100644 --- a/apps/api/app/services/guardrail.py +++ b/apps/api/app/services/guardrail.py @@ -363,6 +363,89 @@ def mask_synthetic_generated_pii(text: str) -> MaskResult: ) +_IDENTITY_SEGMENT_RE = re.compile(r"\s*[·|,/]\s*", re.UNICODE) +_IDENTITY_UUID_RE = re.compile( + r"^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$", + re.IGNORECASE, +) +_IDENTITY_KO_PARTICLE_LOOKAHEAD = ( + r"(?:은|는|이|가|을|를|와|과|의|도|에게|께|랑|하고|님|씨)" +) + + +def _known_identity_variants(identity: str | None) -> list[str]: + """Return conservative display-name variants that are safe to role-tokenize.""" + + raw = str(identity or "").strip() + if not raw or "@" in raw or _IDENTITY_UUID_RE.fullmatch(raw): + return [] + first_segment = _IDENTITY_SEGMENT_RE.split(raw, maxsplit=1)[0].strip() + first_segment = re.sub(r"\s*\(가명\)\s*$", "", first_segment).strip() + variants: list[str] = [] + for candidate in (raw, first_segment): + if candidate in variants: + continue + letters = re.sub(r"[^A-Za-z가-힣]", "", candidate) + if len(letters) < 2: + continue + variants.append(candidate) + return sorted(variants, key=len, reverse=True) + + +def mask_role_identities( + text: str, + *, + counselor_identity: str | None = None, + client_identity: str | None = None, + synthetic_generated: bool = False, +) -> MaskResult: + """Mask known session identities with role tokens, then apply the normal PII gate. + + Only identities already owned by the authenticated session are role-tokenized. + Unknown third-party names keep the generic ``[NAME]`` token so the UI cannot + incorrectly present every person as the client. + """ + + role_values = ( + ("ROLE_COUNSELOR", "[COUNSELOR]", counselor_identity), + ("ROLE_CLIENT", "[CLIENT]", client_identity), + ) + redacted = text + role_entities: list[str] = [] + claimed_variants: set[str] = set() + for entity, placeholder, identity in role_values: + for variant in _known_identity_variants(identity): + normalized = variant.casefold() + if normalized in claimed_variants: + continue + pattern = ( + rf"(?NAME|ORG|PHONE|EMAIL|RRN|NUMID|DATE|MONEY|ADDR)\]" + r"\[(?P