diff --git a/backend/api/core/config.py b/backend/api/core/config.py index 5659ddd3..933812d3 100644 --- a/backend/api/core/config.py +++ b/backend/api/core/config.py @@ -139,7 +139,9 @@ def convert_max_file_size(cls, v): PDF_LINKAGE_MAX_BYTES: int = int( os.getenv('PDF_LINKAGE_MAX_BYTES', str(50 * 1024 * 1024)), ) - PDF_LINKAGE_MAX_RETRIES: int = int(os.getenv('PDF_LINKAGE_MAX_RETRIES', '3')) + PDF_LINKAGE_MAX_RETRIES: int = int( + os.getenv('PDF_LINKAGE_MAX_RETRIES', '3'), + ) PDF_LINKAGE_MAX_CONCURRENCY: int = int( os.getenv('PDF_LINKAGE_MAX_CONCURRENCY', '4'), ) @@ -154,6 +156,21 @@ def convert_max_file_size(cls, v): ) DEBUG: bool = os.getenv('DEBUG', 'false').lower() == 'true' + # Multi-reviewer foundation. Keep disabled until schema and compatibility + # verification has passed in the target environment. + ENABLE_MULTI_REVIEWER_SCHEMA: bool = os.getenv( + 'ENABLE_MULTI_REVIEWER_SCHEMA', 'false', + ).lower().strip() == 'true' + MULTI_REVIEWER_SHADOW_MODE: bool = os.getenv( + 'MULTI_REVIEWER_SHADOW_MODE', 'false', + ).lower().strip() == 'true' + # Schema-changing work is performed by the deployment migration command. + # This is intentionally false in production; it is available for local + # development when an explicit startup bootstrap is desired. + AUTO_MIGRATE: bool = os.getenv( + 'AUTO_MIGRATE', 'true', + ).lower().strip() == 'true' + # ------------------------------------------------------------------------- # Postgres configuration # ------------------------------------------------------------------------- diff --git a/backend/api/migrations/__init__.py b/backend/api/migrations/__init__.py new file mode 100644 index 00000000..c089627f --- /dev/null +++ b/backend/api/migrations/__init__.py @@ -0,0 +1,2 @@ +"""Explicit database migration commands for CAN-SR.""" +from __future__ import annotations diff --git a/backend/api/migrations/__main__.py b/backend/api/migrations/__main__.py new file mode 100644 index 00000000..862119a1 --- /dev/null +++ b/backend/api/migrations/__main__.py @@ -0,0 +1,5 @@ +from __future__ import annotations + +from .cli import main + +main() diff --git a/backend/api/migrations/cli.py b/backend/api/migrations/cli.py new file mode 100644 index 00000000..432c7ead --- /dev/null +++ b/backend/api/migrations/cli.py @@ -0,0 +1,21 @@ +from __future__ import annotations + +import argparse + +from api.core.config import settings +from api.services.review_schema_service import review_schema_service + + +def main() -> None: + parser = argparse.ArgumentParser(description='CAN-SR database migrations') + parser.add_argument('command', choices=['migrate', 'verify', 'status']) + args = parser.parse_args() + if args.command == 'migrate': + result = review_schema_service.migrate(settings.VERSION) + else: + result = review_schema_service.verify_schema() + print(result) + + +if __name__ == '__main__': + main() diff --git a/backend/api/screen/router.py b/backend/api/screen/router.py index 7694defb..7c6b4dce 100644 --- a/backend/api/screen/router.py +++ b/backend/api/screen/router.py @@ -29,6 +29,11 @@ from ..services.cit_db_service import cits_dp_service from ..services.cit_db_service import snake_case from ..services.cit_db_service import snake_case_column +from ..services.postgres_auth import postgres_server +from ..services.review_service import reviewer_identity +from ..services.review_service import ReviewerIdentity +from ..services.review_service import ReviewRepository +from ..services.review_service import WorkUnit from ..services.screening_eligibility_service import compute_screening_decisions from ..services.screening_eligibility_service import screening_eligibility_service from ..services.sr_db_service import srdb_service @@ -1639,6 +1644,33 @@ async def validate_screening_step( detail='Citation not found to update', ) + # Feature-gated canonical dual-write. The legacy write above remains the + # compatibility projection and continues to define the response contract. + if settings.ENABLE_MULTI_REVIEWER_SCHEMA: + canonical_step = 'extract' if step == 'parameters' else step + try: + identity = reviewer_identity(current_user) + unit = WorkUnit( + sr_id=sr_id, + stage=canonical_step, + source_table_name=table_name, + citation_id=citation_id, + criteria_revision=int( + (_sr or {}).get('criteria_revision') or 1, + ), + ) + repository = ReviewRepository(postgres_server.conn) + if checked: + repository.validate(unit, identity, {'checked': True}) + else: + repository.remove_current(unit, identity) + except Exception as e: + logger.exception('Canonical multi-reviewer dual-write failed') + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail=f'Canonical review storage unavailable: {e}', + ) + return { 'status': 'success', 'sr_id': sr_id, diff --git a/backend/api/services/review_schema_service.py b/backend/api/services/review_schema_service.py new file mode 100644 index 00000000..02bac51a --- /dev/null +++ b/backend/api/services/review_schema_service.py @@ -0,0 +1,162 @@ +"""Generic migration runner and read-only schema verifier. + +Migrations are applied explicitly by deployment (or the CLI), while the API +startup path only verifies that the database is compatible with the code. +""" +from __future__ import annotations + +import hashlib +from pathlib import Path +from typing import Any + +MIGRATIONS_PATH = Path(__file__).resolve().parents[2] / 'migrations' +MIGRATION_TABLE = 'review_schema_migrations' +# Compatibility aliases retained for callers/tests from the first bootstrap. +MIGRATION_PATH = MIGRATIONS_PATH / '001_multi_reviewer_schema.sql' +MIGRATION_VERSION = MIGRATION_PATH.stem +REQUIRED_TABLES = { + 'review_schema_migrations', 'review_assignment_policies', + 'review_assignments', 'review_validations', 'reconciliation_cases', + 'reconciliation_decisions', +} + + +def migration_checksum(sql: str | None = None) -> str: + text = sql if sql is not None else '' + return hashlib.sha256(text.encode('utf-8')).hexdigest() + + +def migration_files() -> list[Path]: + return sorted(MIGRATIONS_PATH.glob('[0-9][0-9][0-9]_*.sql')) + + +class ReviewSchemaService: + """Apply and verify the additive schema using an injected DB connection.""" + + def __init__(self, connection_provider=None): + self.connection_provider = connection_provider + + def _provider(self): + if self.connection_provider is None: + # Import lazily so disabled/shadow-mode tests and application + # startup do not require database credentials just to import the + # service module. + from .postgres_auth import postgres_server + self.connection_provider = postgres_server + return self.connection_provider + + def _ensure_migration_table(self, cur) -> None: + cur.execute( + """CREATE TABLE IF NOT EXISTS review_schema_migrations ( + version TEXT PRIMARY KEY, + applied_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + checksum TEXT NOT NULL, + app_version TEXT, + execution_ms INTEGER + )""", + ) + # The first opt-in bootstrap created this table without execution_ms. + # Keep upgrades additive so existing installations can adopt the + # generic runner without a manual repair step. + cur.execute( + 'ALTER TABLE review_schema_migrations ' + 'ADD COLUMN IF NOT EXISTS execution_ms INTEGER', + ) + + def migrate(self, app_version: str | None = None) -> dict[str, Any]: + """Apply all pending migrations in filename order. + + Each migration is committed independently. A failed migration is + rolled back and later migrations are not attempted. + """ + conn = self._provider().conn + files = migration_files() + if not files: + raise RuntimeError(f'No migrations found in {MIGRATIONS_PATH}') + applied: list[str] = [] + for path in files: + version = path.stem + sql = path.read_text(encoding='utf-8') + checksum = hashlib.sha256(sql.encode('utf-8')).hexdigest() + cur = conn.cursor() + try: + cur.execute( + 'SELECT pg_advisory_xact_lock(hashtext(%s))', + ('can_sr_schema_migrations',), + ) + self._ensure_migration_table(cur) + cur.execute( + f'SELECT checksum FROM {MIGRATION_TABLE} WHERE version = %s', + (version,), + ) + existing = cur.fetchone() + if existing: + if existing[0] != checksum: + raise RuntimeError( + f'{version} checksum mismatch; refusing to continue', + ) + conn.commit() + continue + cur.execute(sql) + cur.execute( + f'''INSERT INTO {MIGRATION_TABLE} + (version, checksum, app_version) VALUES (%s, %s, %s)''', + (version, checksum, app_version), + ) + conn.commit() + applied.append(version) + except Exception: + conn.rollback() + raise + finally: + cur.close() + return {'applied': applied, **self.verify_schema(connection=conn)} + + # Backward-compatible name for callers of the original opt-in bootstrap. + def ensure_schema(self, app_version: str | None = None) -> dict[str, Any]: + return self.migrate(app_version) + + def verify_schema(self, connection=None) -> dict[str, Any]: + conn = connection or self._provider().conn + cur = conn.cursor() + try: + cur.execute( + """SELECT table_name FROM information_schema.tables + WHERE table_schema = 'public' + AND table_name = ANY(%s)""", + (sorted(REQUIRED_TABLES),), + ) + tables = {row[0] for row in cur.fetchall()} + missing = sorted(REQUIRED_TABLES - tables) + if missing: + raise RuntimeError( + f'Multi-reviewer schema is incomplete: {missing}', + ) + cur.execute( + f'SELECT version FROM {MIGRATION_TABLE} ORDER BY version', + ) + versions = [row[0] for row in cur.fetchall()] + expected_versions = [path.stem for path in migration_files()] + pending = [ + version for version in expected_versions if version not in versions + ] + if pending: + raise RuntimeError(f'Database migrations pending: {pending}') + # Older installations may retain the first bootstrap marker + # (`multi_reviewer_v1`). It remains useful audit history, but must + # not be treated as a newer migration than numeric file versions. + legacy_versions = [ + version for version in versions if version not in expected_versions + ] + return { + 'version': expected_versions[-1] if expected_versions else None, + 'versions': versions, + 'legacy_versions': legacy_versions, + 'pending': pending, + 'tables': sorted(tables & REQUIRED_TABLES), + } + finally: + cur.close() + + +review_schema_service = ReviewSchemaService() diff --git a/backend/api/services/review_service.py b/backend/api/services/review_service.py new file mode 100644 index 00000000..8694b3ab --- /dev/null +++ b/backend/api/services/review_service.py @@ -0,0 +1,308 @@ +"""Core multi-reviewer domain services. + +This module is deliberately independent of the assignment UI. It provides +stable work-unit identity, authenticated reviewer identity normalization, the +canonical validation repository boundary, and legacy projection helpers. The +existing screening/extraction routers can adopt these operations incrementally +without moving their current response contracts. +""" +from __future__ import annotations + +import json +import re +import unicodedata +import uuid +from dataclasses import dataclass +from datetime import datetime +from datetime import timezone +from typing import Any +from typing import Protocol + + +class ReviewConnection(Protocol): + def cursor(self): ... + + def commit(self) -> None: ... + + def rollback(self) -> None: ... + + +@dataclass(frozen=True) +class WorkUnit: + sr_id: str + stage: str + source_table_name: str + citation_id: int + criteria_revision: int + criterion_key: str | None = None + parameter_key: str | None = None + + def __post_init__(self) -> None: + if self.stage not in {'l1', 'l2', 'extract'}: + raise ValueError(f'Unsupported review stage: {self.stage}') + if self.citation_id < 0: + raise ValueError('citation_id must be non-negative') + if self.parameter_key and not self.criterion_key: + raise ValueError('parameter_key requires criterion_key') + if not self.source_table_name or not re.fullmatch( + r'[A-Za-z_][A-Za-z0-9_]*', self.source_table_name, + ): + raise ValueError('source_table_name is not a safe SQL identifier') + + def values(self) -> tuple[Any, ...]: + return ( + self.sr_id, self.stage, self.source_table_name, self.citation_id, + self.criterion_key, self.parameter_key, self.criteria_revision, + ) + + +@dataclass(frozen=True) +class ReviewerIdentity: + reviewer_id: str + email: str | None = None + + +def reviewer_identity(current_user: dict[str, Any]) -> ReviewerIdentity: + """Derive a stable identity from the authenticated user, never a payload.""" + raw_id = current_user.get('id') or current_user.get( + 'user_id', + ) or current_user.get('sub') + email = current_user.get('email') + if not raw_id and not email: + raise ValueError('Authenticated user has no stable reviewer identity') + normalized_email = str(email).strip().casefold() if email else None + # Legacy tokens may only contain an email. Keep this deterministic until a + # membership service supplies a permanent user ID. + return ReviewerIdentity(str(raw_id or normalized_email), normalized_email) + + +def normalize_extraction_answer(value: Any) -> str: + """Normalize presentation-only differences for agreement checks.""" + text = unicodedata.normalize('NFKC', '' if value is None else str(value)) + text = re.sub(r'<[^>]*>', '', text) + return ' '.join(text.casefold().split()) + + +def extraction_answers_agree(left: Any, right: Any) -> bool: + return normalize_extraction_answer(left) == normalize_extraction_answer(right) + + +class ReviewRepository: + """Persistence boundary for canonical review records. + + SQL is kept here so routers remain responsible for HTTP/authentication and + the legacy projector remains independently testable. + """ + + def __init__(self, connection: ReviewConnection): + self.connection = connection + + @staticmethod + def _work_unit_where(unit: WorkUnit) -> tuple[str, tuple[Any, ...]]: + return ( + '''sr_id = %s AND stage = %s AND source_table_name = %s + AND citation_id = %s AND criteria_revision = %s + AND criterion_key IS NOT DISTINCT FROM %s + AND parameter_key IS NOT DISTINCT FROM %s''', + ( + unit.sr_id, unit.stage, unit.source_table_name, unit.citation_id, + unit.criteria_revision, unit.criterion_key, unit.parameter_key, + ), + ) + + def save_draft( + self, + unit: WorkUnit, + reviewer: ReviewerIdentity, + answer: Any, + *, + explanation: str | None = None, + evidence: Any = None, + assignment_id: str | None = None, + client_request_id: str | None = None, + expected_version: int | None = None, + ) -> dict[str, Any]: + return self._save( + unit, reviewer, answer, 'draft', explanation, evidence, + assignment_id, client_request_id, expected_version, + ) + + def validate( + self, + unit: WorkUnit, + reviewer: ReviewerIdentity, + answer: Any, + *, + explanation: str | None = None, + evidence: Any = None, + assignment_id: str | None = None, + client_request_id: str | None = None, + expected_version: int | None = None, + ) -> dict[str, Any]: + return self._save( + unit, reviewer, answer, 'validated', explanation, evidence, + assignment_id, client_request_id, expected_version, + ) + + def visible_validations( + self, unit: WorkUnit, reviewer: ReviewerIdentity, + ) -> list[dict[str, Any]]: + """Return only records the current reviewer is allowed to see. + + Drafts are private. Other reviewers' validated answers become visible + only after the current reviewer has locked their own validation. + """ + conn = self.connection + cur = conn.cursor() + try: + where, params = self._work_unit_where(unit) + cur.execute( + f'''SELECT id, reviewer_id, state, revision, answer_json, + explanation, evidence_json, created_at, validated_at, version + FROM review_validations + WHERE {where} AND state IN ('draft', 'validated') + ORDER BY reviewer_id, revision''', + params, + ) + rows = cur.fetchall() + own_validated = any( + row[1] == reviewer.reviewer_id and row[2] == 'validated' + for row in rows + ) + visible = [ + row for row in rows + if row[1] == reviewer.reviewer_id + or (own_validated and row[2] == 'validated') + ] + return [ + { + 'id': row[0], 'reviewer_id': row[1], 'state': row[2], + 'revision': row[3], 'answer': row[4], + 'explanation': row[5], 'evidence': row[6], + 'created_at': row[7], 'validated_at': row[8], + 'version': row[9], + } + for row in visible + ] + finally: + cur.close() + + def remove_current(self, unit: WorkUnit, reviewer: ReviewerIdentity) -> int: + """Remove the current draft/validated record for a reviewer.""" + cur = self.connection.cursor() + try: + where, params = self._work_unit_where(unit) + cur.execute( + f'''DELETE FROM review_validations + WHERE {where} AND reviewer_id = %s + AND state IN ('draft', 'validated')''', + (*params, reviewer.reviewer_id), + ) + deleted = cur.rowcount + self.connection.commit() + return deleted + except Exception: + self.connection.rollback() + raise + finally: + cur.close() + + def _save( + self, unit: WorkUnit, reviewer: ReviewerIdentity, answer: Any, + state: str, explanation: str | None, evidence: Any, + assignment_id: str | None, client_request_id: str | None, + expected_version: int | None, + ) -> dict[str, Any]: + if state not in {'draft', 'validated'}: + raise ValueError(f'Unsupported write state: {state}') + conn = self.connection + cur = conn.cursor() + row_id = str(uuid.uuid4()) + try: + where, params = self._work_unit_where(unit) + if client_request_id: + cur.execute( + '''SELECT id, reviewer_id, state, revision, version + FROM review_validations + WHERE reviewer_id = %s AND client_request_id = %s''', + (reviewer.reviewer_id, client_request_id), + ) + request_row = cur.fetchone() + if request_row: + # Retrying a request returns the original result and does + # not create another revision or overwrite another user. + conn.rollback() + return { + 'id': request_row[0], 'reviewer_id': request_row[1], + 'state': request_row[2], 'revision': request_row[3], + 'version': request_row[4], 'idempotent_replay': True, + } + cur.execute( + f'''SELECT id, version, state FROM review_validations + WHERE {where} AND reviewer_id = %s + AND state IN ('draft', 'validated') + ORDER BY revision DESC LIMIT 1 FOR UPDATE''', + (*params, reviewer.reviewer_id), + ) + current = cur.fetchone() + if current and expected_version is not None and current[1] != expected_version: + raise RuntimeError('review validation version conflict') + revision = 1 + if current: + cur.execute( + 'SELECT COALESCE(MAX(revision), 0) + 1 FROM review_validations ' + f'WHERE {where} AND reviewer_id = %s', + (*params, reviewer.reviewer_id), + ) + revision = int(cur.fetchone()[0]) + cur.execute( + "UPDATE review_validations SET state = 'superseded', version = version + 1 " + 'WHERE id = %s', (current[0],), + ) + now = datetime.now(timezone.utc) + cur.execute( + '''INSERT INTO review_validations + (id, assignment_id, sr_id, stage, source_table_name, citation_id, + criterion_key, parameter_key, criteria_revision, reviewer_id, + revision, state, answer_json, explanation, evidence_json, source, + client_request_id, created_at, validated_at) + VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, + %s::jsonb, %s, %s::jsonb, 'user', %s, %s, %s)''', + ( + row_id, assignment_id, unit.sr_id, unit.stage, + unit.source_table_name, unit.citation_id, unit.criterion_key, + unit.parameter_key, unit.criteria_revision, reviewer.reviewer_id, + revision, state, json.dumps(answer), explanation, + json.dumps(evidence) if evidence is not None else None, + client_request_id, now, now if state == 'validated' else None, + ), + ) + conn.commit() + return { + 'id': row_id, 'reviewer_id': reviewer.reviewer_id, + 'state': state, 'revision': revision, 'version': 1, + } + except Exception: + conn.rollback() + raise + finally: + cur.close() + + +class CompatibilityProjector: + """Pure mapping helpers for legacy response/column projections.""" + + @staticmethod + def validation_entry(record: dict[str, Any]) -> dict[str, str]: + reviewer = record.get('reviewer_id') or record.get('email') or '' + validated_at = record.get( + 'validated_at', + ) or record.get('created_at') or '' + return {'user': str(reviewer), 'validated_at': str(validated_at)} + + @staticmethod + def legacy_validation_list(records: list[dict[str, Any]]) -> list[dict[str, str]]: + return [ + CompatibilityProjector.validation_entry(record) + for record in records if record.get('state') == 'validated' + ] diff --git a/backend/deploy.sh b/backend/deploy.sh index aac3e25a..af7f838c 100755 --- a/backend/deploy.sh +++ b/backend/deploy.sh @@ -158,6 +158,17 @@ sleep 10 # Start main API echo -e "${BLUE}๐ Starting main API...${NC}" +# Read the feature flag from the compose environment (.env), not only from +# the shell running this script. The migration job is deliberately explicit; +# the API container itself only verifies the schema at startup. +docker compose run --rm --no-deps api sh -lc '\ + if [ "$$ENABLE_MULTI_REVIEWER_SCHEMA" = "true" ]; then \ + echo "๐งฑ Applying database migrations..."; \ + python -m api.migrations migrate; \ + echo "โ Database migrations applied"; \ + else \ + echo "โน๏ธ Multi-reviewer migrations disabled"; \ + fi' if [ "$DEV" = true ]; then docker compose up -d api --remove-orphans else diff --git a/backend/main.py b/backend/main.py index 645910b3..fd465707 100644 --- a/backend/main.py +++ b/backend/main.py @@ -40,6 +40,25 @@ async def startup_event(): ) print('๐ Starting CAN-SR Backend...', flush=True) + if settings.ENABLE_MULTI_REVIEWER_SCHEMA: + try: + from api.services.review_schema_service import review_schema_service + print('๐ฅ Verifying multi-reviewer database schema...', flush=True) + if settings.AUTO_MIGRATE: + print( + 'โ ๏ธ AUTO_MIGRATE is enabled; applying pending migrations', flush=True, + ) + result = await run_in_threadpool(review_schema_service.migrate, settings.VERSION) + else: + result = await run_in_threadpool(review_schema_service.verify_schema) + print( + f"โ Multi-reviewer schema verified ({result['version']})", flush=True, + ) + except Exception as e: + # Opt-in foundation bootstrap must fail loudly when enabled; the + # default-disabled path remains identical to the current app. + print(f'โ Multi-reviewer schema bootstrap failed: {e}', flush=True) + raise print('๐ Initializing systematic review database...', flush=True) # Ensure systematic review table exists in PostgreSQL try: @@ -136,7 +155,9 @@ async def _stale_job_reaper(): from api.services.fulltext_attachment_service import reconcile_pending_blob_cleanup cleaned = await reconcile_pending_blob_cleanup() if cleaned: - print(f'๐งน Removed {cleaned} orphaned full-text blobs', flush=True) + print( + f'๐งน Removed {cleaned} orphaned full-text blobs', flush=True, + ) except Exception as reaper_err: print( f"โ ๏ธ Stale job reaper error: {reaper_err}", flush=True, diff --git a/backend/migrations/001_multi_reviewer_schema.sql b/backend/migrations/001_multi_reviewer_schema.sql new file mode 100644 index 00000000..771906ce --- /dev/null +++ b/backend/migrations/001_multi_reviewer_schema.sql @@ -0,0 +1,155 @@ +-- CAN-SR multi-reviewer foundation schema. +-- This migration is additive: legacy citation tables and human_* columns are +-- intentionally not altered. + +CREATE TABLE IF NOT EXISTS review_assignment_policies ( + id UUID PRIMARY KEY, + sr_id TEXT NOT NULL, + stage TEXT NOT NULL CHECK (stage IN ('l1', 'l2', 'extract')), + default_granularity TEXT NOT NULL + CHECK (default_granularity IN ('citation', 'criterion', 'parameter')), + overlap SMALLINT NOT NULL CHECK (overlap BETWEEN 1 AND 3), + eligible_reviewer_ids JSONB NOT NULL DEFAULT '[]'::jsonb, + workload_targets JSONB NOT NULL DEFAULT '{}'::jsonb, + visibility_rule TEXT NOT NULL DEFAULT 'after_own_validation', + agreement_rule_version TEXT NOT NULL DEFAULT 'v1', + criteria_revision INTEGER NOT NULL DEFAULT 1, + due_at TIMESTAMPTZ, + version BIGINT NOT NULL DEFAULT 1, + created_by TEXT NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + UNIQUE (sr_id, stage, version) +); + +-- Policy versions are unique. The repository resolves the greatest version +-- for an SR/stage; PostgreSQL partial-index predicates cannot contain a +-- subquery, so โactiveโ is intentionally a service-level concept. + +CREATE TABLE IF NOT EXISTS review_assignments ( + id UUID PRIMARY KEY, + policy_id UUID NOT NULL REFERENCES review_assignment_policies(id), + sr_id TEXT NOT NULL, + stage TEXT NOT NULL CHECK (stage IN ('l1', 'l2', 'extract')), + source_table_name TEXT NOT NULL, + citation_id BIGINT NOT NULL, + criterion_key TEXT, + parameter_key TEXT, + criteria_revision INTEGER NOT NULL, + reviewer_id TEXT NOT NULL, + slot SMALLINT NOT NULL CHECK (slot >= 1), + source TEXT NOT NULL CHECK (source IN + ('stage_default', 'bulk_rule', 'citation', 'criterion', 'parameter')), + status TEXT NOT NULL CHECK (status IN ('active', 'removed', 'stale')), + priority INTEGER NOT NULL DEFAULT 0, + due_at TIMESTAMPTZ, + version BIGINT NOT NULL DEFAULT 1, + created_by TEXT NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + CHECK (parameter_key IS NULL OR criterion_key IS NOT NULL) +); + +CREATE UNIQUE INDEX IF NOT EXISTS uq_review_active_assignment_reviewer + ON review_assignments ( + sr_id, stage, source_table_name, citation_id, + COALESCE(criterion_key, ''), COALESCE(parameter_key, ''), reviewer_id + ) WHERE status = 'active'; + +CREATE UNIQUE INDEX IF NOT EXISTS uq_review_active_assignment_slot + ON review_assignments ( + sr_id, stage, source_table_name, citation_id, + COALESCE(criterion_key, ''), COALESCE(parameter_key, ''), slot + ) WHERE status = 'active'; + +CREATE INDEX IF NOT EXISTS ix_review_assignment_reviewer + ON review_assignments (sr_id, stage, reviewer_id, status); +CREATE INDEX IF NOT EXISTS ix_review_assignment_work_unit + ON review_assignments (sr_id, stage, source_table_name, citation_id, + criterion_key, parameter_key); +CREATE INDEX IF NOT EXISTS ix_review_assignment_queue + ON review_assignments (sr_id, stage, status, due_at); + +CREATE TABLE IF NOT EXISTS review_validations ( + id UUID PRIMARY KEY, + assignment_id UUID REFERENCES review_assignments(id), + sr_id TEXT NOT NULL, + stage TEXT NOT NULL CHECK (stage IN ('l1', 'l2', 'extract')), + source_table_name TEXT NOT NULL, + citation_id BIGINT NOT NULL, + criterion_key TEXT, + parameter_key TEXT, + criteria_revision INTEGER NOT NULL, + reviewer_id TEXT NOT NULL, + revision INTEGER NOT NULL CHECK (revision >= 1), + state TEXT NOT NULL CHECK (state IN + ('draft', 'validated', 'returned', 'superseded', 'stale')), + answer_json JSONB NOT NULL, + explanation TEXT, + evidence_json JSONB, + source TEXT NOT NULL CHECK (source IN ('user', 'legacy_migration')), + client_request_id TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + validated_at TIMESTAMPTZ, + returned_at TIMESTAMPTZ, + superseded_by UUID REFERENCES review_validations(id), + version BIGINT NOT NULL DEFAULT 1, + CHECK (parameter_key IS NULL OR criterion_key IS NOT NULL) +); + +CREATE UNIQUE INDEX IF NOT EXISTS uq_review_validation_request + ON review_validations (reviewer_id, client_request_id) + WHERE client_request_id IS NOT NULL; +CREATE UNIQUE INDEX IF NOT EXISTS uq_review_current_validation + ON review_validations ( + sr_id, stage, source_table_name, citation_id, + COALESCE(criterion_key, ''), COALESCE(parameter_key, ''), + reviewer_id, state + ) WHERE state IN ('draft', 'validated'); +CREATE INDEX IF NOT EXISTS ix_review_validation_reviewer + ON review_validations (sr_id, stage, reviewer_id, state, validated_at); +CREATE INDEX IF NOT EXISTS ix_review_validation_work_unit + ON review_validations (sr_id, stage, source_table_name, citation_id, + criterion_key, parameter_key, state); + +CREATE TABLE IF NOT EXISTS reconciliation_cases ( + id UUID PRIMARY KEY, + sr_id TEXT NOT NULL, + stage TEXT NOT NULL CHECK (stage IN ('l1', 'l2', 'extract')), + source_table_name TEXT NOT NULL, + citation_id BIGINT NOT NULL, + criterion_key TEXT, + parameter_key TEXT, + case_version INTEGER NOT NULL CHECK (case_version >= 1), + status TEXT NOT NULL CHECK (status IN + ('pending', 'in_progress', 'resolved', 'reopened')), + reason TEXT NOT NULL, + detected_from_validation_ids JSONB NOT NULL DEFAULT '[]'::jsonb, + assigned_reviewer_id TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + resolved_at TIMESTAMPTZ, + version BIGINT NOT NULL DEFAULT 1, + CHECK (parameter_key IS NULL OR criterion_key IS NOT NULL) +); + +CREATE UNIQUE INDEX IF NOT EXISTS uq_review_active_reconciliation + ON reconciliation_cases ( + sr_id, stage, source_table_name, citation_id, + COALESCE(criterion_key, ''), COALESCE(parameter_key, '') + ) WHERE status IN ('pending', 'in_progress', 'reopened'); + +CREATE TABLE IF NOT EXISTS reconciliation_decisions ( + id UUID PRIMARY KEY, + case_id UUID NOT NULL REFERENCES reconciliation_cases(id), + case_version INTEGER NOT NULL, + decided_by TEXT NOT NULL, + final_answer_json JSONB NOT NULL, + rationale TEXT NOT NULL, + evidence_json JSONB, + preferred_validation_id UUID REFERENCES review_validations(id), + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS ix_reconciliation_queue + ON reconciliation_cases (sr_id, stage, status, updated_at); diff --git a/backend/tests/review_schema_service_test.py b/backend/tests/review_schema_service_test.py new file mode 100644 index 00000000..ccd84b73 --- /dev/null +++ b/backend/tests/review_schema_service_test.py @@ -0,0 +1,36 @@ +from __future__ import annotations + +from pathlib import Path + +from api.services import review_schema_service as schema + + +def test_migration_file_is_resolved_inside_backend(): + assert schema.MIGRATION_PATH == Path(__file__).resolve( + ).parents[1] / 'migrations' / '001_multi_reviewer_schema.sql' + assert schema.MIGRATION_PATH.is_file() + + +def test_migration_files_are_ordered(): + files = schema.migration_files() + assert files == sorted(files) + assert files[0].stem == '001_multi_reviewer_schema' + + +def test_canonical_version_is_based_on_migration_files_not_legacy_markers(): + assert schema.MIGRATION_VERSION == '001_multi_reviewer_schema' + + +def test_migration_checksum_is_deterministic(): + sql = 'CREATE TABLE foundation (id integer);' + assert schema.migration_checksum(sql) == schema.migration_checksum(sql) + assert len(schema.migration_checksum(sql)) == 64 + + +def test_feature_is_disabled_by_default(monkeypatch): + # The imported settings object is intentionally not mutated. This checks + # the source-level default contract without requiring a database. + config = Path(__file__).resolve().parents[1] / 'api' / 'core' / 'config.py' + source = config.read_text(encoding='utf-8') + assert "ENABLE_MULTI_REVIEWER_SCHEMA', 'false'" in source + assert "AUTO_MIGRATE', 'false'" in source diff --git a/backend/tests/review_service_test.py b/backend/tests/review_service_test.py new file mode 100644 index 00000000..b162aea4 --- /dev/null +++ b/backend/tests/review_service_test.py @@ -0,0 +1,74 @@ +from __future__ import annotations + +import pytest +from api.services.review_service import CompatibilityProjector +from api.services.review_service import extraction_answers_agree +from api.services.review_service import normalize_extraction_answer +from api.services.review_service import reviewer_identity +from api.services.review_service import WorkUnit + + +def test_work_unit_is_stable_and_parameter_requires_criterion(): + unit = WorkUnit( + sr_id='sr-1', stage='extract', source_table_name='citations_1', + citation_id=42, criteria_revision=3, criterion_key='pico', + parameter_key='population', + ) + assert unit.values() == ( + 'sr-1', 'extract', 'citations_1', 42, 'pico', 'population', 3, + ) + + with pytest.raises(ValueError, match='requires criterion'): + WorkUnit( + sr_id='sr-1', stage='extract', source_table_name='citations_1', + citation_id=42, criteria_revision=3, parameter_key='population', + ) + + +def test_work_unit_rejects_unsafe_dynamic_table_name(): + with pytest.raises(ValueError, match='safe SQL identifier'): + WorkUnit( + sr_id='sr-1', stage='l1', source_table_name='citations; DROP TABLE x', + citation_id=1, criteria_revision=1, + ) + + +def test_reviewer_identity_comes_from_authenticated_user_and_normalizes_email(): + identity = reviewer_identity( + {'id': 'user-7', 'email': ' Reviewer@Example.COM '}, + ) + assert identity.reviewer_id == 'user-7' + assert identity.email == 'reviewer@example.com' + assert reviewer_identity( + {'email': 'Reviewer@Example.COM'}, + ).reviewer_id == 'reviewer@example.com' + + with pytest.raises(ValueError, match='stable reviewer identity'): + reviewer_identity({}) + + +def test_extraction_agreement_only_ignores_presentation_differences(): + assert normalize_extraction_answer(' A\nB ') == 'a b' + assert extraction_answers_agree('
A
', 'a') + assert not extraction_answers_agree('10 mg', '10') + assert not extraction_answers_agree( + 'heart attack', 'myocardial infarction', + ) + + +def test_legacy_projection_contains_only_validated_records(): + records = [ + {'reviewer_id': 'u1', 'state': 'draft', 'created_at': 'draft-time'}, + {'reviewer_id': 'u1', 'state': 'validated', 'validated_at': 't1'}, + {'reviewer_id': 'u2', 'state': 'validated', 'validated_at': 't2'}, + ] + assert CompatibilityProjector.legacy_validation_list(records) == [ + {'user': 'u1', 'validated_at': 't1'}, + {'user': 'u2', 'validated_at': 't2'}, + ] + + +def test_projection_preserves_legacy_user_and_timestamp_shape(): + assert CompatibilityProjector.validation_entry( + {'email': 'reviewer@example.com', 'created_at': 't0'}, + ) == {'user': 'reviewer@example.com', 'validated_at': 't0'}