"""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