"""Preprocess sub-phase 3 — data validation: independent integrity + graded coherence checks. Each check is an INDEPENDENT named boolean column with a ``chk_`` prefix; there is no tiering between them — they differ only in how they're consumed: integrity checks (binary) — logically impossible facts; ANY hit flips ``validation_failed`` (the candidate is treated as untrustworthy, not merely low-quality). graded checks (weighted) — "smells off but possible" signals; each contributes a small penalty that sums into a capped ``incoherence_score`` used softly downstream. Runs as one vectorized pass on RAW input (before enrich) so flags reflect what the candidate submitted. Flags only — NEVER drops a row; the hard filter that acts on these is a later phase. """ from __future__ import annotations from datetime import date import polars as pl from common.logging import step from .load import to_date # Columns this stage appends to every row (the downstream contract). Order is display-only. FLAG_COLUMNS = [ "validation_failed", "validation_rules", # integrity verdict + audit trail "incoherence_score", "coherence_rules", # graded score + audit trail "years_exp_claimed", "years_exp_computed", # claimed vs duration-computed YoE, side by side "verified_skill_count", "trust_score", # quality / identity-verification counts "date_order_ok", "salary_range_ok", "active_validated", "experience_validated", # 0/1 markers "chk_expert_zero", "chk_career_dates_invalid", "chk_yoe_impossible", # integrity checks "chk_date_order_invalid", "chk_salary_incoherent", # graded checks ] # Max gap (years) between claimed and computed YoE for the claim to count as "validated" — 6 months. EXPERIENCE_TOLERANCE_YEARS = 0.5 # Upper bound on summed graded penalties — keeps one noisy profile from dominating scoring. COHERENCE_CAP = 0.5 # Integrity checks (binary): any True => the row is hard-failed. _INTEGRITY_CHECKS = ["chk_expert_zero", "chk_career_dates_invalid", "chk_yoe_impossible"] # Graded checks: (check_column, penalty weight). Heavier weight = stronger incoherence signal. _GRADED_CHECKS = [("chk_date_order_invalid", 0.15), ("chk_salary_incoherent", 0.10)] def _derived_exprs(as_of: date) -> list[pl.Expr]: """Intermediate ``d_*`` columns the checks depend on (career-history dates only). Skill durations are intentionally NOT summed: skills run in parallel within a role, so summed skill-months aren't commensurable with calendar career-months. Only career dates give a sound bound. """ field = pl.element().struct.field earliest_start = pl.col("career_history").list.eval(to_date(field("start_date"))).list.min() return [ # Outer bound on plausible experience; claimed YoE shouldn't exceed this by much. ((pl.lit(as_of) - earliest_start).dt.total_days() / 365.25) .cast(pl.Float32).alias("d_career_span_years"), # YoE implied by summing role durations (independent of the self-reported field). (pl.col("career_history").list.eval(field("duration_months")).list.sum() / 12.0) .cast(pl.Float32).alias("d_computed_yoe"), ] def _check_exprs() -> list[pl.Expr]: """The three integrity checks plus the date-order graded signal (booleans).""" field = pl.element().struct.field return [ # Claims advanced/expert mastery yet logged zero months on the skill — contradictory. pl.col("skills").list.eval( field("proficiency").is_in(["advanced", "expert"]) & (field("duration_months") == 0) ).list.any().alias("chk_expert_zero"), # A role that ends before it starts is impossible; null dates count as fine (not a hit). pl.col("career_history").list.eval( (to_date(field("end_date")) < to_date(field("start_date"))).fill_null(False) # noqa: E501 ).list.any().alias("chk_career_dates_invalid"), # Last-active before signup is a graded timeline glitch rather than a disqualifier. (to_date(pl.col("last_active_date")) < to_date(pl.col("signup_date"))) .fill_null(False).alias("chk_date_order_invalid"), # Claimed YoE exceeds the calendar span by >2yr — impossible (2yr slack absorbs gaps/overlap). (pl.col("years_of_experience") > pl.col("d_career_span_years") + 2) .fill_null(False).alias("chk_yoe_impossible"), ] def _coherence_exprs() -> list[pl.Expr]: """Boolean "smells off but possible" columns; each maps to a penalty in ``_GRADED_CHECKS``.""" return [ # Inverted salary band, or a max so large it's plainly a unit/typo error (>1000 in k-units). ( (pl.col("expected_salary_min") > pl.col("expected_salary_max")) | (pl.col("expected_salary_max") > 1000) ).fill_null(False).alias("chk_salary_incoherent"), ] def _audit_trail(checks: list[str], alias: str) -> pl.Expr: """Space-joined names of whichever of ``checks`` fired on a row — a human-readable audit string.""" fired = [pl.when(pl.col(check)).then(pl.lit(check)).otherwise(pl.lit("")) for check in checks] return pl.concat_str(fired, separator=" ").str.strip_chars().alias(alias) def _assemble_exprs() -> list[pl.Expr]: """Fold the raw check columns into the user-facing flags, scores, and audit-trail strings.""" # Hard-fail if ANY integrity check fired. failed = pl.any_horizontal([pl.col(check) for check in _INTEGRITY_CHECKS]).alias("validation_failed") # Sum each fired graded check's penalty weight, then clamp at COHERENCE_CAP. weighted_penalty = sum( (pl.col(check).cast(pl.Float64) * weight for check, weight in _GRADED_CHECKS), pl.lit(0.0) ) incoherence = ( pl.min_horizontal(pl.lit(COHERENCE_CAP), weighted_penalty).cast(pl.Float32).alias("incoherence_score") ) # Identity-verification tally (email + phone + linkedin), 0..3. trust = ( pl.col("verified_email").cast(pl.Int8) + pl.col("verified_phone").cast(pl.Int8) + pl.col("linkedin_connected").cast(pl.Int8) ).alias("trust_score") return [ failed, _audit_trail(_INTEGRITY_CHECKS, "validation_rules"), incoherence, _audit_trail([check for check, _ in _GRADED_CHECKS], "coherence_rules"), # Count of skills carrying an assessment-score entry — a quality proxy. pl.col("verified_skills").list.len().cast(pl.Int8).alias("verified_skill_count"), trust, pl.col("years_of_experience").cast(pl.Float32).alias("years_exp_claimed"), pl.col("d_computed_yoe").cast(pl.Float32).alias("years_exp_computed"), ] def validate(candidates: pl.DataFrame, as_of: date) -> pl.DataFrame: """Vectorized integrity + graded coherence checks; adds the check + flag/score columns. Never drops.""" with step("preprocess.validate", rows_in=candidates.height) as metrics: # Staged so each stage can reference the previous: derived d_* → checks → coherence → assemble → 0/1. validated = ( candidates.with_columns(_derived_exprs(as_of)) .with_columns(_check_exprs()) .with_columns(_coherence_exprs()) .with_columns(_assemble_exprs()) .with_columns( date_order_ok=(~pl.col("chk_date_order_invalid")).cast(pl.Int8), salary_range_ok=(~pl.col("chk_salary_incoherent")).cast(pl.Int8), ) .with_columns( # 1 iff signup precedes/equals last-active; missing dates fall to 0 (stricter than the marker). active_validated=(to_date(pl.col("signup_date")) <= to_date(pl.col("last_active_date"))) .fill_null(False).cast(pl.Int8), # 1 iff claimed YoE agrees with the duration-summed computed YoE to within 6 months. experience_validated=( (pl.col("years_of_experience") - pl.col("d_computed_yoe")).abs() <= EXPERIENCE_TOLERANCE_YEARS ).fill_null(False).cast(pl.Int8), ) ) failed_count = int(validated.select(pl.col("validation_failed").sum()).item()) metrics["rows_out"] = validated.height metrics["validation_failed"] = failed_count metrics["invalid_rate"] = round(failed_count / validated.height, 5) if validated.height else 0.0 return validated