File size: 10,845 Bytes
700e896
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
f85f907
700e896
fb57bfa
700e896
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
a41b023
 
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
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
"""Preprocess sub-phase β€” materialize: one vectorized pass to features + the embedded text body.

Lean scope (this build): structural counts (from the raw ``flat_skills`` etc.), ``experience_band``
computed inline from years-of-experience, sentinel-clean platform signals, the salary band, and the
deterministic templated ``text`` that retrieval embeds (identity + experience + skills + education β€” NO
behavioral/platform signals) plus its content hash. The 7 quality pillars + ``candidate_quality`` and the
rerank-feature legs are deferred to their own later stages (they are ranking-side features).
"""

from __future__ import annotations

import hashlib
from datetime import date

import polars as pl

from common.logging import step

from .load import to_date

# Proficiency levels that count as a "core" (strong) skill.
_PROFICIENT = ["advanced", "expert"]
# Notice period (days) under which a candidate counts as "available soon".
_AVAILABLE_SOON_DAYS = 30
# Float inputs left at source precision β€” every other derived float is rounded to 3 dp in one pass.
_RAW_FLOAT_PASSTHROUGH = {
    "years_of_experience", "profile_completeness_score", "recruiter_response_rate",
    "avg_response_time_hours", "expected_salary_min", "expected_salary_max",
    "github_activity_score", "interview_completion_rate", "offer_acceptance_rate",
}


def _text(column: str) -> pl.Expr:
    """Null-safe string view of a column (null β†’ empty)."""
    return pl.col(column).cast(pl.String).fill_null("")


def _experience_band_expr() -> pl.Expr:
    """Left-closed seniority band from years-of-experience (computed inline, no enrich)."""
    yoe = pl.col("years_of_experience")
    return (
        pl.when(yoe < 3).then(pl.lit("junior"))
        .when(yoe < 6).then(pl.lit("mid"))
        .when(yoe < 9).then(pl.lit("senior"))
        .when(yoe < 13).then(pl.lit("lead"))
        .otherwise(pl.lit("principal"))
        .alias("experience_band")
    )


def _count_exprs() -> list[pl.Expr]:
    """Structural counts over the nested arrays + raw flat lists (skills via raw ``flat_skills``)."""
    field = pl.element().struct.field
    durations = pl.col("career_history").list.eval(field("duration_months").cast(pl.Int64).fill_null(0))
    assessed = pl.col("skill_assessment_scores").list.eval(
        pl.when(field("score") >= 0).then(field("score")).otherwise(None)
    )
    advanced = pl.col("skills").list.eval(field("proficiency").is_in(_PROFICIENT).cast(pl.Int32)).list.sum()
    expert = pl.col("skills").list.eval((field("proficiency") == "expert").cast(pl.Int32)).list.sum()
    num_skills = pl.col("flat_skills").list.len()
    return [
        num_skills.alias("num_skills"),
        pl.col("career_history").list.len().alias("num_roles"),
        durations.list.mean().fill_null(0.0).alias("avg_tenure_months"),
        pl.col("career_history").list.eval(field("company").cast(pl.String)).list.n_unique().alias("num_companies"),
        pl.col("flat_career_industries").list.eval(pl.element().filter(pl.element() != "")).list.n_unique().alias("num_industries"),
        advanced.alias("num_advanced_skills"),
        expert.alias("num_expert_skills"),
        assessed.list.len().alias("num_assessments_taken"),
        assessed.list.mean().alias("avg_skill_assessment"),
        assessed.list.max().alias("max_skill_assessment"),
        pl.when(num_skills > 0).then(assessed.list.len() / num_skills).otherwise(None).alias("assessment_coverage"),
        pl.col("education").list.len().alias("num_degrees"),
        pl.col("certifications").list.len().alias("num_certifications"),
        (
            pl.col("verified_email").cast(pl.Int32) + pl.col("verified_phone").cast(pl.Int32)
            + pl.col("linkedin_connected").cast(pl.Int32)
        ).alias("verification_count"),
        (pl.col("notice_period_days") <= _AVAILABLE_SOON_DAYS).alias("is_available_soon"),
        pl.col("career_history").list.eval(
            pl.when(field("is_current")).then(field("duration_months").cast(pl.Int64)).otherwise(None)
        ).list.max().fill_null(0).alias("current_role_tenure_months"),
    ]


def _signal_exprs(as_of: date) -> list[pl.Expr]:
    """Date deltas + sentinel-aware clean signals (negative github/offer == 'no data')."""
    github = pl.col("github_activity_score")
    offer = pl.col("offer_acceptance_rate")
    salary_min = pl.col("expected_salary_min").fill_null(0.0)
    salary_max = pl.col("expected_salary_max").fill_null(0.0)
    return [
        ((pl.lit(as_of) - to_date(pl.col("last_active_date"))).dt.total_days().fill_null(999).cast(pl.Int32)).alias("days_since_active"),
        ((pl.lit(as_of) - to_date(pl.col("signup_date"))).dt.total_days().fill_null(0).cast(pl.Int32)).alias("account_age_days"),
        (github >= 0).alias("has_github"),
        pl.when(github >= 0).then(github.clip(0, 100)).otherwise(None).alias("github_score_clean"),
        (offer >= 0).alias("has_offer_history"),
        pl.when(offer >= 0).then(offer.clip(0, 1)).otherwise(None).alias("offer_acceptance_rate_clean"),
        salary_min.alias("salary_min"),
        salary_max.alias("salary_max"),
        ((salary_min + salary_max) / 2.0).alias("salary_mid"),
    ]


def _text_section_exprs() -> list[pl.Expr]:
    """Per-section text blocks fed into the embedded body (identity / skills / experience / education)."""
    field = pl.element().struct.field
    industries = pl.col("flat_career_industries").list.eval(pl.element().filter(pl.element() != "")).list.unique()
    identity = pl.concat_str(
        [_text("headline"), _text("summary"), _text("current_title"), pl.lit(" @ "), _text("current_company")],
        separator="\n",
    ).str.strip_chars()
    core = pl.col("skills").list.eval(
        pl.when(field("proficiency").is_in(_PROFICIENT)).then(field("name").cast(pl.String)).otherwise(None)
    ).list.drop_nulls().list.join(", ")
    familiar = pl.col("skills").list.eval(
        pl.when(field("proficiency").is_in(["beginner", "intermediate"])).then(field("name").cast(pl.String)).otherwise(None)
    ).list.drop_nulls().list.join(", ")
    experience = pl.col("career_history").list.eval(
        pl.concat_str([
            field("title").cast(pl.String).fill_null(""), pl.lit(" @ "),
            field("company").cast(pl.String).fill_null(""), pl.lit(" ("),
            field("industry").cast(pl.String).fill_null(""), pl.lit(") - "),
            field("description").cast(pl.String).fill_null(""),
        ])
    ).list.head(8).list.join("\n")
    education = pl.col("education").list.eval(
        pl.concat_str([
            field("degree").cast(pl.String).fill_null(""), pl.lit(" in "),
            field("field_of_study").cast(pl.String).fill_null(""), pl.lit(", "),
            field("institution").cast(pl.String).fill_null(""),
        ])
    ).list.join("\n")
    return [
        industries.list.join(", ").alias("industries_text"),
        identity.alias("identity_text"),
        core.alias("core_skills_text"),
        familiar.alias("familiar_skills_text"),
        experience.alias("experience_text"),
        education.alias("education_text"),
        pl.col("flat_certifications").list.eval(pl.element().filter(pl.element() != "")).list.join(", ").alias("certifications_text"),
    ]


def _summary_text_expr() -> pl.Expr:
    """The deterministic, templated candidate summary β€” THE embedded body (no platform signals)."""
    years = pl.col("years_of_experience").round(0).cast(pl.Int64).cast(pl.String)
    best_tier = pl.col("flat_institution_tiers").list.eval(pl.element().filter(pl.element() != "")).list.max().fill_null("")
    lead = pl.concat_str([
        _text("current_title"), pl.lit(" with "), years, pl.lit(" years of experience"),
        pl.when(pl.col("industries_text").str.len_chars() > 0).then(pl.lit(" in ") + pl.col("industries_text")).otherwise(pl.lit("")),
        pl.lit(" ("), _text("experience_band"), pl.lit("-band)."),
    ])
    education_clause = pl.when(pl.col("num_degrees") > 0).then(
        pl.lit("Education: ") + pl.col("education_text").str.replace_all("\n", "; ")
        + pl.when(best_tier.str.len_chars() > 0).then(pl.lit(" (") + best_tier + pl.lit(" tier).")).otherwise(pl.lit("."))
    ).otherwise(pl.lit(""))
    return pl.concat_str(
        [
            lead,
            _clause("Currently at ", _text("current_company")),
            _clause("Location: ", _text("location")),  # so cities are visible to bm25 + the embedders
            _clause("Core skills: ", pl.col("core_skills_text")),
            _clause("Also familiar with: ", pl.col("familiar_skills_text")),
            _clause("Recent roles: ", pl.col("experience_text").str.replace_all("\n", "; ")),
            education_clause,
            _clause("Certifications: ", pl.col("certifications_text")),
        ],
        separator=" ", ignore_nulls=True,
    ).str.replace_all(r"\s+", " ").str.strip_chars()


def _clause(label: str, value: pl.Expr) -> pl.Expr:
    """Prefix ``value`` with ``label`` when non-empty, else null (dropped by concat ignore_nulls)."""
    return pl.when(value.str.len_chars() > 0).then(pl.lit(label) + value).otherwise(None)


def _sha16(values: list[str]) -> list[str]:
    """sha1[:16] of each summary body β€” changes when the embedded text changes (re-embed trigger)."""
    return [hashlib.sha1(value.encode("utf-8")).hexdigest()[:16] for value in values]


def materialize(candidates: pl.DataFrame, as_of: date) -> pl.DataFrame:
    """Add structural counts, clean signals, experience_band, and the embedded ``text`` body. Never drops."""
    with step("preprocess.materialize", rows_in=candidates.height) as metrics:
        # Ordered passes β€” text sections must exist before the summary that composes them.
        out = candidates.with_columns(
            _experience_band_expr(), *_count_exprs(), *_signal_exprs(as_of)
        ).with_columns(_text_section_exprs())
        out = out.with_columns(_summary_text_expr().alias("summary_text"))
        out = out.with_columns(pl.col("summary_text").alias("text"))
        # Provisional hash β€” the next preprocess stage (quality.py) appends a strengths clause to ``text`` and
        # re-hashes it there, so this value is superseded before anything reads it in the normal pipeline order.
        out = out.with_columns(
            pl.Series("embed_source_hash", _sha16(out.get_column("summary_text").to_list()))
        )
        # Single rounding pass: every derived float to 3 dp, leaving raw source signals untouched.
        round_cols = [
            name for name, dtype in out.schema.items()
            if dtype in (pl.Float32, pl.Float64) and name not in _RAW_FLOAT_PASSTHROUGH
        ]
        out = out.with_columns([pl.col(name).round(3) for name in round_cols])
        metrics["rows_out"] = out.height
        metrics["columns"] = out.width
        return out