HR-ATS-Portal/backend/job/candidate/views.py

453 lines
20 KiB
Python

from sqlalchemy.ext.asyncio import AsyncSession
import asyncio,dataclasses,hashlib,os,logging,io,uuid
from datetime import datetime,timezone
from fastapi import HTTPException
from pypdf import PdfReader
from sqlalchemy import select
from sqlalchemy.orm import selectinload
from sqlmodel import true
from app.core.errors import ATSError,ErrorCode
from app.models.scoring import CompletedCandidate
from app.services.pdf import extract_resume,sanitize_filename
from app.services.scoring import score_batch
from inbox.models import Inbox_Messages,Inbox
from job.candidate.models import Candidates
from job.candidate.plugins import (
FILE_NOT_FOUND,
build_job_description,
get_scorer,
get_scoring_settings,
normalize_spaced_text,
)
from job.candidate.serializers import serialize_candidate,serialize_candidate_profile
from job.job_post.models import JobPosts
from job.job_post.serializers import serialize_job_post
from job.candidate.models import Notes
from job.notes.serializers import serialize_note
class FileRead:
def __init__(self,session:AsyncSession,filename=None,file=None):
self.session=session
self.filename=filename
self.file=file
async def read_file(self,file=None,filename=None):
try:
reader = PdfReader(io.BytesIO(self.file))
if reader.is_encrypted:
raise HTTPException(400, "PDF is password protected")
pages = [(page.extract_text() or "") for page in reader.pages]
return {
"filename": self.filename,
"num_pages": len(reader.pages),
"text": normalize_spaced_text("\n".join(pages)),
}
except HTTPException:
raise
except Exception as e:
raise HTTPException(400, str(e))
async def match_inbox_cv(self,inbox_message_id):
from inbox.plugins import resolve_attachment_path
from inbox.tasks import match_inbox_message
row=await Inbox_Messages.get_inbox_message_by_id(self.session,inbox_message_id)
if not row:
raise HTTPException(status_code=404,detail="Message not found")
if not row.attachment or not row.file_path:
raise HTTPException(status_code=400,detail="your file isnt in the system")
found=None
for path_str in (p.strip() for p in row.file_path.split(",") if p.strip()):
path=resolve_attachment_path(path_str)
if path.is_file():
found=path
break
if found is None:
raise HTTPException(status_code=400,detail="your file isnt in the system")
created_at=datetime.now(timezone.utc).isoformat()
task=await match_inbox_message.kicker().with_labels(
created_at=created_at,
correlation_id=str(row.id),
queue="inbox",
).kiq(str(row.id),force=True)
file_name=(row.file_name or "").split(",")[0].strip() or found.name
return {
"queued":True,
"inbox_message_id":str(row.id),
"file_name":file_name,
"task_id":task.task_id,
}
# async def get_intention(self,input):
# try:
# get_subject=Inbox_Messages.candidate_x_inbox(self.session,self.candidate_id)
# get_file=
class CandidateScoring:
"""ATS scoring of CVs against one job post, persisted to the candidates table.
Complements the agent's inbox-match flow: the agent suggests WHICH job a CV is
for; this service scores HOW WELL a CV fits a chosen job (0-100 leaderboard).
Deviation from the bulk-ats HTTP API (which rejects a whole batch with 415/413
on a bad file): here per-file problems become persisted rows with
status="failed" so one broken attachment never sinks the rest of the batch.
Request-level errors (unknown job, too many files) still raise.
"""
def __init__(self,session:AsyncSession):
self.session=session
async def score_uploads(self,job_id,files,current_user):
"""files: list of (filename, bytes) pairs from the route handler."""
settings=get_scoring_settings()
if len(files)>settings.max_resumes_per_request:
raise HTTPException(
status_code=413,
detail=f"At most {settings.max_resumes_per_request} resumes per request",
)
sources=[]
for filename,data in files:
source={
"filename":filename or "resume.pdf",
"data":data,
"file_path":None,
"inbox_message_id":None,
"precheck":None,
}
if not (filename or "").lower().endswith(".pdf"):
source["precheck"]=(ErrorCode.UNSUPPORTED_FILE_TYPE,"Only PDF resumes are supported.")
elif len(data)>settings.max_pdf_size_bytes:
source["precheck"]=(ErrorCode.PAYLOAD_TOO_LARGE,"The file exceeds the size limit.")
sources.append(source)
return await self._score_and_persist(job_id,sources,"upload",current_user)
async def score_inbox(self,job_id,message_ids,current_user):
"""Score the decoded attachments of inbox messages (PK uuids, not Graph ids)."""
# Local import: inbox.plugins imports this module (FileRead), so a top-level
# import would be circular — same pattern as match_inbox_cv above.
from inbox.plugins import resolve_attachment_path
sources=[]
for mid in message_ids:
row=await Inbox_Messages.get_inbox_message_by_id(self.session,mid)
if row is None:
raise HTTPException(status_code=404,detail=f"Inbox message {mid} not found")
if not row.file_path:
continue
for path_str in (p.strip() for p in row.file_path.split(",") if p.strip()):
path=resolve_attachment_path(path_str)
source={
"filename":path.name,
"data":None,
"file_path":str(path),
"inbox_message_id":row.id,
"precheck":None,
}
suffix=path.suffix.lower()
if suffix in (".doc",".docx"):
source["precheck"]=(ErrorCode.UNSUPPORTED_FILE_TYPE,"DOC/DOCX extraction is not supported yet.")
elif suffix!=".pdf":
source["precheck"]=(ErrorCode.UNSUPPORTED_FILE_TYPE,"Only PDF resumes are supported.")
elif not path.is_file():
source["precheck"]=(FILE_NOT_FOUND,"The decoded attachment is missing on disk.")
else:
try:
source["data"]=await asyncio.to_thread(path.read_bytes)
except OSError:
source["precheck"]=(FILE_NOT_FOUND,"The decoded attachment could not be read.")
sources.append(source)
if not sources:
raise HTTPException(status_code=400,detail="No attachments found for the given message(s)")
return await self._score_and_persist(job_id,sources,"inbox",current_user)
async def fetch_candidates(self,job_id=None):
# job_id omitted -> the whole pool across jobs (frontend Candidates/TalentPool).
if job_id is not None:
job=await JobPosts.get_job_post_by_id(self.session,job_id)
if job is None or job.is_deleted:
raise HTTPException(status_code=404,detail="Job post not found")
rows=await Candidates.get_candidates_by_job(self.session,job_id)
return [serialize_candidate(row) for row in rows]
async def fetch_candidate_by_id(self,candidate_id):
row=await Candidates.get_candidate_by_id(self.session,candidate_id)
if row is None:
raise HTTPException(status_code=404,detail="Candidate not found")
return serialize_candidate(row)
async def _score_and_persist(self,job_id,sources,source_kind,current_user):
job=await JobPosts.get_job_post_by_id(self.session,job_id)
if job is None or job.is_deleted:
raise HTTPException(status_code=404,detail="Job post not found")
settings=get_scoring_settings()
jd=build_job_description(job)
if len(jd)>settings.max_jd_chars:
raise HTTPException(status_code=422,detail="The job post is too large to score against")
# Slot-indexed like app/api/routes.py: results merge back by position, never
# by filename — inbox attachments can share a basename.
results_by_slot={}
extracted=[]
for slot,source in enumerate(sources):
source["safe_name"]=sanitize_filename(source["filename"])
data=source["data"]
source["sha256"]=hashlib.sha256(data).hexdigest() if data is not None else None
if source["precheck"] is not None:
code,message=source["precheck"]
results_by_slot[slot]=self._failed_fields(source,code,message)
continue
try:
# pypdf is CPU-bound: keep it off the event loop. Despace BEFORE
# scoring so keyword verification sees the exact text the model saw;
# ExtractedResume is frozen, hence dataclasses.replace.
resume=await asyncio.to_thread(
extract_resume,data,source["safe_name"],settings.max_resume_chars
)
resume=dataclasses.replace(resume,text=normalize_spaced_text(resume.text))
except ATSError as exc:
results_by_slot[slot]=self._failed_fields(source,exc.error_code,exc.public_message)
continue
extracted.append((slot,resume))
scored=await score_batch(
[resume for _,resume in extracted],
job_description=jd,
scorer=get_scorer(),
concurrency=settings.scoring_concurrency,
)
for (slot,_),result in zip(extracted,scored,strict=True):
source=sources[slot]
if isinstance(result,CompletedCandidate):
results_by_slot[slot]={
**self._base_fields(source),
"status":"completed",
"candidate_name":result.candidate_name,
"job_title":result.job_title,
"current_company":result.current_company,
"years_experience":result.years_experience,
"match_score":result.match_score,
"matched_keywords":result.matched_keywords,
"missing_keywords":result.missing_keywords,
"summary_critique":result.summary_critique,
"error_code":None,
"error_message":None,
}
else:
results_by_slot[slot]=self._failed_fields(source,result.error_code,result.error_message)
common={
"job_id":job.id,
"source":source_kind,
"created_by":uuid.UUID(str(current_user["id"])),
"model":settings.openai_model,
}
rows=[]
for slot in range(len(sources)):
fields={**results_by_slot[slot],**common}
rows.append(await Candidates.upsert_candidate(self.session,fields))
# Leaderboard order: completed by score desc, failures last, stable.
rows.sort(key=lambda r:(0,-(r.match_score or 0)) if r.status=="completed" else (1,0))
return [serialize_candidate(row) for row in rows]
@staticmethod
def _base_fields(source):
return {
"inbox_message_id":source["inbox_message_id"],
"filename":source["safe_name"],
"file_path":source["file_path"],
"content_sha256":source["sha256"],
}
@classmethod
def _failed_fields(cls,source,code,message):
return {
**cls._base_fields(source),
"status":"failed",
"error_code":str(code),
"error_message":message,
"candidate_name":None,
"job_title":None,
"current_company":None,
"years_experience":None,
"match_score":None,
"matched_keywords":[],
"missing_keywords":[],
"summary_critique":None,
}
class CandidateView:
def __init__(self,session:AsyncSession):
self.session=session
async def get_candidate(self,user_id=None,limit=10,offset=0,search=None):
try:
detail=bool(user_id)
# Detail mode must see every application for the candidate, not one page.
fetch_limit=1000 if detail else limit
rows=await Inbox.get_candidate_profile(session=self.session,user_id=user_id,limit=fetch_limit,offset=offset,search=search)
if detail:
return await self.attach_profile_detail(rows)
return await self.attach_job_posts(rows)
except HTTPException:
raise
except Exception as e:
raise HTTPException(status_code=500,detail=str(e))
async def count_candidates(self,user_id=None,search=None):
try:
return await Inbox.count_candidate_profiles(session=self.session,user_id=user_id,search=search)
except HTTPException:
raise
except Exception as e:
raise HTTPException(status_code=500,detail=str(e))
async def update_candidate(self,user_id,payload):
try:
if not user_id:
raise HTTPException(status_code=400,detail="user_id is required")
fields={k:v for k,v in (payload or {}).items() if k in ("favorite","rating") and v is not None}
if not fields:
raise HTTPException(status_code=400,detail="favorite or rating is required")
links=await Inbox.get_candidate_profile(session=self.session,user_id=user_id,limit=100,offset=0)
records=links if isinstance(links,list) else ([links] if links else [])
if not records:
raise HTTPException(status_code=404,detail="Candidate not found")
for link in records:
await Inbox.update_inbox(self.session,link.id,fields)
refreshed=await Inbox.get_candidate_profile(session=self.session,user_id=user_id,limit=100,offset=0)
return await self.attach_profile_detail(refreshed)
except HTTPException:
raise
except Exception as e:
raise HTTPException(status_code=500,detail=str(e))
async def get_job_post_by_id(self,record_id,data=None,*,as_assigned=False):
"""Load full job_posts row and optionally append it onto a candidate payload."""
try:
job_post_data=await JobPosts.get_job_post_by_id(session=self.session,record_id=record_id)
if not job_post_data:
return None
payload=serialize_job_post(job_post_data)
if isinstance(data,dict):
if as_assigned:
data["assigned_job_post"]=payload
if payload.get("created_by_name"):
data["recruiter"]=payload.get("created_by_name")
data["recruiter_id"]=payload.get("created_by")
if payload.get("title"):
data["job_title"]=payload.get("title")
else:
data.setdefault("job_posts",[]).append(payload)
if data.get("recruiter") is None and payload.get("created_by_name"):
data["recruiter"]=payload.get("created_by_name")
data["recruiter_id"]=payload.get("created_by")
if data.get("job_title") is None and payload.get("title"):
data["job_title"]=payload.get("title")
return payload
except Exception as e:
raise HTTPException(status_code=500,detail=str(e))
async def attach_job_posts(self,data):
"""Normalize list/single, serialize each record, attach full job_posts rows."""
single=not isinstance(data,list)
records=[data] if single else list(data or [])
enriched=[]
for record in records:
payload=serialize_candidate_profile(record)
payload["job_posts"]=[]
payload["assigned_job_post"]=None
assigned_id=payload.get("assigned_job_post_id")
if assigned_id:
await self.get_job_post_by_id(record_id=assigned_id,data=payload,as_assigned=True)
for job_id in payload.get("suggested_job_post_ids") or []:
await self.get_job_post_by_id(record_id=job_id,data=payload)
enriched.append(payload)
return enriched[0] if single else enriched
async def attach_profile_detail(self,data):
"""Detail mode: flatten child collections across every Inbox row for the candidate."""
single=not isinstance(data,list)
records=[data] if single else list(data or [])
if not records:
return {} if single else []
interviews=[]
activity=[]
feedback=[]
documents=[]
job_posts=[]
assigned_job_post=None
base=None
user_id=None
favorite=None
rating=None
for record in records:
payload=serialize_candidate_profile(record,detail=True)
if base is None:
base=payload
user_id=payload.get("user_id")
favorite=payload.get("favorite")
rating=payload.get("rating")
interviews.extend(payload.get("interviews") or [])
activity.extend(payload.get("activity") or [])
feedback.extend(payload.get("feedback") or [])
documents.extend(payload.get("documents") or [])
if payload.get("assigned_job_post_id") and assigned_job_post is None:
await self.get_job_post_by_id(
record_id=payload.get("assigned_job_post_id"),
data=payload,
as_assigned=True,
)
assigned_job_post=payload.get("assigned_job_post")
if base.get("recruiter") is None and payload.get("recruiter"):
base["recruiter"]=payload.get("recruiter")
base["recruiter_id"]=payload.get("recruiter_id")
if base.get("job_title") is None and payload.get("job_title"):
base["job_title"]=payload.get("job_title")
for job_id in payload.get("suggested_job_post_ids") or []:
await self.get_job_post_by_id(record_id=job_id,data=payload)
for jp in payload.get("job_posts") or []:
if not any(x.get("id")==jp.get("id") for x in job_posts):
job_posts.append(jp)
if base.get("recruiter") is None and payload.get("recruiter"):
base["recruiter"]=payload.get("recruiter")
base["recruiter_id"]=payload.get("recruiter_id")
if base.get("job_title") is None and payload.get("job_title"):
base["job_title"]=payload.get("job_title")
notes=[]
uid=Notes._as_uuid(user_id) if user_id else None
if uid is not None:
result=await self.session.execute(
select(Notes)
.options(selectinload(Notes.author))
.where(Notes.user_id==uid)
.order_by(Notes.created_at.desc())
)
notes=[serialize_note(r) for r in result.scalars().all()]
activity.sort(key=lambda r:(r.get("activity_date") or ""),reverse=True)
base["favorite"]=favorite
base["rating"]=rating
base["interviews"]=interviews
base["activity"]=activity
base["feedback"]=feedback
base["documents"]=documents
base["notes"]=notes
base["job_posts"]=job_posts or base.get("job_posts") or []
base["assigned_job_post"]=assigned_job_post
if assigned_job_post:
base["assigned_job_post_id"]=assigned_job_post.get("id")
if assigned_job_post.get("created_by_name"):
base["recruiter"]=assigned_job_post.get("created_by_name")
base["recruiter_id"]=assigned_job_post.get("created_by")
if assigned_job_post.get("title"):
base["job_title"]=assigned_job_post.get("title")
return base