File size: 8,466 Bytes
700e896
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
"""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