from sqlalchemy.ext.asyncio import AsyncSession import asyncio,base64,dataclasses,hashlib,io,logging,os,uuid from datetime import datetime,timezone from pathlib import Path from dotenv import load_dotenv from fastapi import HTTPException from pypdf import PdfReader 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,AtsResults from job.candidate.models import Candidates from job.candidate.plugins import ( FILE_NOT_FOUND, build_job_description, candidate_completed_fields, candidate_failed_fields, contained_download_path, documents_from_message, extract_pdf_link_uris, get_scorer, get_scoring_settings, normalize_spaced_text, ) from job.candidate.serializers import serialize_candidate,serialize_candidate_profile,serialize_manual_candidate_profile,serialize_manual_upload_candidate,serialize_matching_candidate,serialize_manager_candidate from job.job_post.models import JobPosts from job.job_post.serializers import serialize_job_post from job.candidate.models import Notes,Manual_UPLOAD_CANDIDATE from job.history.enums import HistoryEvent from job.history.views import HistoryRecorder from job.notes.serializers import serialize_note from job.candidate.plugins import extract_candidate_email from users.models import Users from users.permissions import is_hiring_manager,sees_all_candidates from employment_agent.plugins import parse_phone load_dotenv() logger=logging.getLogger("job.candidate.views") CV_QUEUE_NAME=os.getenv("TASKIQ_CV_QUEUE_NAME","cv_upload") MANUAL_UPLOAD_TO_ADDRESS=os.getenv( "MANUAL_UPLOAD_TO_ADDRESS","manual-cv-upload@hr-ats.local" ) MANAGER_SCOPE_DETAIL="You can only access candidates allocated to jobs opened from your requisitions" CREATOR_SCOPE_DETAIL="You can only access candidates allocated to jobs you created" async def assigned_job_ids_for_user(session,user_id): """Job posts this candidate is allocated to (inbox assignment + manual upload).""" ids=set() if not user_id: return ids rows=await Inbox.get_candidate_profile(session=session,user_id=user_id,limit=1000,offset=0) records=rows if isinstance(rows,list) else ([rows] if rows else []) for rec in records: msg=getattr(rec,"messages",None) jid=getattr(msg,"assigned_job_post_id",None) if msg is not None else None if jid: ids.add(jid) manual=await Manual_UPLOAD_CANDIDATE.get_by_user_id(session,user_id) if manual and manual.job_post_id: ids.add(manual.job_post_id) return ids async def job_id_for_application(session,inbox_id=None,manual_id=None): if inbox_id is not None: link=await Inbox.get_inbox_with_message(session,inbox_id) if link is None: return None,None msg=link.messages return (msg.assigned_job_post_id if msg is not None else None),link.user_id if manual_id is not None: manual=await Manual_UPLOAD_CANDIDATE.get_by_id(session,manual_id) if manual is None: return None,None return manual.job_post_id,manual.user_id return None,None async def owned_job_ids_for_candidate_scope(session,current_user): """Job ids this user may see, or None when the list is unscoped. Hiring manager → requisition / assigned-manager jobs. candidates.manage or admin → None (all applications). Otherwise → job_posts.created_by = this user. Never role_id. """ if is_hiring_manager(current_user): return await JobPosts.ids_for_manager(session,current_user.get("id")) if sees_all_candidates(current_user): return None return await JobPosts.ids_for_creator(session,current_user.get("id")) async def job_post_ids_for_candidate_list(session,current_user,assigned_job_post_id=None): """None = unscoped list. [] = nothing visible. Else UUID list for the query.""" owned=await owned_job_ids_for_candidate_scope(session,current_user) requested=JobPosts._as_uuid(assigned_job_post_id) if assigned_job_post_id is not None else None if owned is None: return [requested] if requested else None if requested is not None: return [requested] if requested in set(owned) else [] return list(owned) def _scope_detail(current_user): if is_hiring_manager(current_user): return MANAGER_SCOPE_DETAIL return CREATOR_SCOPE_DETAIL async def assert_manager_candidate_access( session,current_user,*,user_id=None,job_post_id=None,inbox_id=None,manual_id=None, ): """Row access: hiring-manager jobs, unscoped (manage/admin), or jobs this user created.""" if not is_hiring_manager(current_user) and sees_all_candidates(current_user): return owned=await owned_job_ids_for_candidate_scope(session,current_user) owned=set(owned or []) detail=_scope_detail(current_user) if not owned: raise HTTPException(status_code=403,detail=detail) job_id=JobPosts._as_uuid(job_post_id) if job_post_id is not None else None uid=user_id if job_id is None and (inbox_id is not None or manual_id is not None): job_id,uid=await job_id_for_application(session,inbox_id=inbox_id,manual_id=manual_id) if job_id is None and uid is not None: candidate_jobs=await assigned_job_ids_for_user(session,uid) if candidate_jobs & owned: return raise HTTPException(status_code=403,detail=detail) if job_id is None or job_id not in owned: raise HTTPException(status_code=403,detail=detail) async def parse_linkedin_url_from_cv(resume_text) -> str | None: """Employment-agent `linkedin_url` from CV text. None if absent, sentinel, invented, or the call fails.""" text=(resume_text or "").strip() if not text: return None try: from employment_agent.execute_agent import run_employment_agent from employment_agent.plugins import parse_linkedin fields=await run_employment_agent(resume_text=text) url_fields=parse_linkedin({"linkedin_url":fields.get("linkedin_url") or ""},text) return url_fields.get("linkedin_url") except Exception: logger.exception("employment agent linkedin_url parse failed") return None 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] text = normalize_spaced_text("\n".join(pages)) # Icon-only LinkedIn buttons never appear in extract_text(); the # URL is on the annotation. Append so the employment agent can # return linkedin_url as its own parsed key. uris = extract_pdf_link_uris(reader) if uris: extra = "\n".join(uris) text = f"{text}\n\n{extra}".strip() if text else extra return { "filename": self.filename, "num_pages": len(reader.pages), "text": text, } except HTTPException: raise except Exception as e: raise HTTPException(400, str(e)) async def injest_manual_upload(self): try: parsed=await self.read_file() text=(parsed.get("text") or "").strip() if not text: raise HTTPException(status_code=400,detail="No usable text could be extracted from the PDF") return parsed except HTTPException: raise except Exception as e: raise HTTPException(status_code=400,detail=str(e)) async def save_manual_upload(self): """Deprecated — Manual CVs go to S3 via create_candidate (no local disk).""" raise HTTPException( status_code=410, detail="Local CV storage was removed; use create_candidate (S3 Manual/{id}/{user_id}/)", ) @staticmethod def discard_upload(file_path): """Best-effort removal of a leftover local CV (legacy rows only).""" if not file_path: return if str(file_path).lower().startswith("http://") or str(file_path).lower().startswith("https://"): return try: Path(file_path).unlink(missing_ok=True) except OSError as e: logger.warning("could not remove orphaned upload %s: %s",file_path,e) async def ingest_upload(self,candidate_email=None,candidate_name=None,current_user=None): """Persist a recruiter-uploaded CV with full email-ingestion parity (S3).""" from inbox.file_decoder import extract_pdf_attachments from inbox.cv_tasks import match_uploaded_cv from inbox.plugins import attach_email_pdfs_to_s3 from inbox.views import Email from s3.plugins import S3ServiceError,assert_pdf parsed=await self.read_file() text=parsed.get("text") or "" parsed_linkedin=await parse_linkedin_url_from_cv(text) detected,emails_found=extract_candidate_email(text) supplied=(candidate_email or "").strip().lower() or None email=supplied or detected email_source="recruiter" if supplied else ("cv" if detected else None) if not email: raise HTTPException( status_code=422, detail={ "error_code":"CANDIDATE_EMAIL_REQUIRED", "filename":parsed.get("filename"), "num_pages":parsed.get("num_pages"), "emails_found":emails_found, "text":text, }, ) filename=self.filename or "resume.pdf" try: assert_pdf(filename,"application/pdf") except S3ServiceError as e: raise HTTPException(status_code=e.status_code,detail=e.message) from e # In-memory only — no decoded_attachments write. import base64 as _b64 pdfs=extract_pdf_attachments([{ "name":filename, "contentBytes":_b64.b64encode(self.file).decode("ascii"), }]) if not pdfs: raise HTTPException(status_code=400,detail="Only PDF resumes are allowed") now=datetime.now(timezone.utc).isoformat() email_data={ "id":f"manual-cv:{uuid.uuid4()}", "subject":f"Manual CV upload — {filename}", "body":{"content":"","contentType":"text"}, "hasAttachments":True, "attachments":[{"name":filename}], "from":{"emailAddress":{"address":email,"name":(candidate_name or "").strip()}}, "toRecipients":[{"emailAddress":{"address":MANUAL_UPLOAD_TO_ADDRESS}}], "ccRecipients":[], "bccRecipients":[], "replyTo":[], "isRead":False, "sentDateTime":now, "receivedDateTime":now, } row,new_user_email=await Inbox_Messages.insert_email( self.session,email_data,file_path=None, ) try: row=await attach_email_pdfs_to_s3(self.session,row,pdfs,created_new=True) except Exception as e: raise HTTPException(status_code=502,detail=f"S3 upload failed: {e}") from e if parsed_linkedin: user_id=await Inbox_Messages.get_linked_user_id(self.session,row.id) if user_id and await Users.set_linkedin_url_if_empty( self.session,user_id=user_id,url=parsed_linkedin, ): await self.session.commit() created_at=datetime.now(timezone.utc).isoformat() # Enqueue failure must not fail the upload: the CV row is already # persisted, and with the broker down (optional locally) the kiq call # raises a connection error. Suggestions just arrive later, or never. try: task=await match_uploaded_cv.kicker().with_labels( created_at=created_at, correlation_id=str(row.id), queue=CV_QUEUE_NAME, ).kiq(str(row.id),force=False) except Exception as e: logger.warning("cv match enqueue skipped for %s: %s",row.id,e) task=None account_setup=None if new_user_email: try: account_setup=await Email(session=self.session).send_account_setup( [new_user_email] ) except Exception as e: logger.warning("account setup mail failed for %s: %s",new_user_email,e) account_setup=[{"email":new_user_email,"sent":False}] user=await Users.get_user_by_email(self.session,email) if user: await HistoryRecorder(self.session).record( HistoryEvent.CANDIDATE_IMPORTED.value, current_user=current_user,user_id=user.id, entity_type="inbox_message",entity_id=row.id, to_value=email,description=f"CV uploaded: {filename}",commit=True, ) await HistoryRecorder(self.session).record( HistoryEvent.DOCUMENT_UPLOADED.value, current_user=current_user,user_id=user.id, entity_type="document",entity_id=row.id, to_value=filename,commit=True, ) return { "queued":True, "inbox_message_id":str(row.id), "task_id":task.task_id if task else None, "filename":parsed.get("filename"), "num_pages":parsed.get("num_pages"), "candidate_email":email, "email_source":email_source, "account_setup":account_setup, "text":text, } async def match_inbox_cv(self,inbox_message_id,current_user=None): from inbox.plugins import load_file_bytes 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_name=None for path_str in (p.strip() for p in row.file_path.split(",") if p.strip()): # S3 URL or local — presence of bytes (or a https URL we already stored) counts. if path_str.lower().startswith("http://") or path_str.lower().startswith("https://"): found_name=Path(path_str.replace("\\","/")).name or "resume.pdf" break raw=load_file_bytes(path_str) if raw is not None: found_name=Path(path_str.replace("\\","/")).name or "resume.pdf" break if found_name 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 await HistoryRecorder(self.session).record( HistoryEvent.CANDIDATE_IMPORTED.value, current_user=current_user,message_id=inbox_message_id, entity_type="inbox_message",entity_id=row.id, to_value=file_name,description=f"CV uploaded: {file_name}",commit=True, ) 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, "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 PDF attachments of inbox messages (S3 URLs or legacy local paths).""" from inbox.plugins import load_file_bytes 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 names=[n.strip() for n in (row.file_name or "").split(",") if n.strip()] for idx,path_str in enumerate(p.strip() for p in row.file_path.split(",") if p.strip()): name=names[idx] if idx 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,total=await Candidates.get_candidates_by_job( self.session,job_id,limit=limit,offset=offset, ) return [serialize_candidate(row) for row in rows],total 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") fields_by_slot=await self._score_sources(sources,jd,settings) # Prefer the Manual S3 URL for this email+job when scoring from a raw upload # (Add Candidate scores right after create — same link as manual_upload_candidate). for slot,source in enumerate(sources): fields=fields_by_slot.get(slot) or {} if (fields.get("file_path") or source.get("file_path") or "").strip(): continue email=(fields.get("candidate_email") or source.get("candidate_email") or "").strip().lower() if not email: continue manual=await Manual_UPLOAD_CANDIDATE.get_by_email_and_job(self.session,email,job.id) if manual and (manual.file_path or "").strip(): fields["file_path"]=manual.file_path.strip() source["file_path"]=manual.file_path.strip() 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={**fields_by_slot[slot],**common} email=(fields.get("candidate_email") or "").strip().lower() if email and not fields.get("linkedin_url"): user=await Users.get_user_by_email(self.session,email) if user and (user.linkedin_url or "").strip(): fields["linkedin_url"]=user.linkedin_url row=await Candidates.upsert_candidate(self.session,fields) if row.linkedin_url and row.candidate_email: if await Users.set_linkedin_url_if_empty( self.session,email=row.candidate_email,url=row.linkedin_url, ): await self.session.commit() rows.append(row) await self._sync_ats_results(source_kind,job,rows,sources,current_user) 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] async def _score_sources(self,sources,jd,settings): # Slot-indexed: results merge back by position, never by filename — # inbox attachments can share a basename. fields_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"] fields_by_slot[slot]=candidate_failed_fields(source,code,message) continue try: 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)) if not source.get("candidate_email"): detected,_=extract_candidate_email(resume.text) source["candidate_email"]=detected except ATSError as exc: fields_by_slot[slot]=candidate_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,_resume),result in zip(extracted,scored,strict=True): source=sources[slot] if isinstance(result,CompletedCandidate): fields_by_slot[slot]=candidate_completed_fields(source,result) else: fields_by_slot[slot]=candidate_failed_fields(source,result.error_code,result.error_message) return fields_by_slot async def _sync_ats_results(self,source_kind,job,rows,sources,current_user=None): if source_kind=="inbox": best={} for source,row in zip(sources,rows): mid=source.get("inbox_message_id") if row.status=="completed" and mid: cur=best.get(mid) if cur is None or (row.match_score or 0)>(cur.match_score or 0): best[mid]=row for message_id,row in best.items(): try: await self._sync_inbox_ats(message_id,job,row,current_user=current_user) except Exception: await self.session.rollback() logger.exception("inbox ATS denorm failed for message %s",message_id) return for row in rows: if row.status!="completed": continue try: await self._sync_upload_ats(job,row,current_user=current_user) except Exception: await self.session.rollback() logger.exception("upload ATS history failed for candidate %s",row.id) async def _sync_inbox_ats(self,message_id,job,row,current_user=None): """Land a completed score on inbox_messages / inbox / ats_results. message_id is the scoring call's known inbox_messages PK, not a column read off the Candidates row. """ msg=await Inbox_Messages.get_inbox_message_by_id(self.session,message_id) if msg is None: return band=CandidateView._recommendation(row.match_score) or "" denorm=True assigned=msg.assigned_job_post_id link=await Inbox.get_inbox_by_message_id(self.session,message_id) if assigned and str(assigned)!=str(job.id) and link is not None: existing=await AtsResults.get_for_inbox_job(self.session,link.id,assigned) denorm=existing is None if denorm: await Inbox_Messages.set_ats_score(self.session,message_id,row.match_score,band) if link is None: return old=await AtsResults.get_current_for_inbox(self.session,link.id) old_score=old.overall_score if old else None identity=await AtsResults.resolve_identity(self.session,row.candidate_email,row.id) await AtsResults.insert_result(self.session,{ "inbox_id":link.id, **identity, "job_post_id":job.id, "overall_score":float(row.match_score), "band":band, "model_name":row.model, "is_current":True, }) await HistoryRecorder(self.session).record( HistoryEvent.ATS_SCORED.value, current_user=current_user,user_id=link.user_id,inbox_id=link.id, entity_type="ats_result",entity_id=row.id, from_value=old_score,to_value=row.match_score, description=f"{band} against {job.title or 'job'}",commit=True, ) async def _sync_upload_ats(self,job,row,current_user=None): """History row for an upload-sourced score — inbox_id stays NULL.""" band=CandidateView._recommendation(row.match_score) or "" identity=await AtsResults.resolve_identity(self.session,row.candidate_email,row.id) old=None if identity.get("user_id"): old=await AtsResults.get_current_for_user(self.session,identity["user_id"],job.id) elif identity.get("candidate_id"): old=await AtsResults.get_current_for_candidate(self.session,identity["candidate_id"]) old_score=old.overall_score if old else None await AtsResults.insert_result(self.session,{ "inbox_id":None, **identity, "job_post_id":job.id, "overall_score":float(row.match_score), "band":band, "model_name":row.model, "is_current":True, }) await HistoryRecorder(self.session).record( HistoryEvent.ATS_SCORED.value, current_user=current_user,user_id=identity.get("user_id"), entity_type="ats_result",entity_id=row.id, from_value=old_score,to_value=row.match_score, description=f"{band} against {job.title or 'job'}",commit=True, ) class CandidateView: def __init__(self,session:AsyncSession): self.session=session @staticmethod def _recommendation(score): if score is None: return None # Same bands the frontend uses (Candidates.jsx / seed.js). return "Strong Match" if score>=82 else "Potential Match" if score>=65 else "Weak Match" @staticmethod def _score_from_message(record): """Inbox denorm on the already-loaded messages row — not a Candidates join.""" msg=getattr(record,"messages",None) if msg is None or msg.ats_score is None: return None,None band=(msg.ats_band or "").strip() or None return msg.ats_score,band or CandidateView._recommendation(msg.ats_score) async def create_candidate(self,candidate_email=None,candidate_name=None,candidate_phone=None,job_post_id=None,current_company=None,current_position=None,platform=None,experience=None,status=None,referral_by=None,file_name=None,file_path=None,full_text=None,current_user=None,file_bytes=None,content_type=None): """Create manual_upload_candidate, then S3 upload under Manual/{id}/{user_id}/. Atomicity: if S3 fails after the row insert, the row is deleted (rolled back). PDF gate runs before any DB write when file_bytes is supplied. """ from s3.plugins import S3,S3ServiceError,S3Source,assert_pdf row=None try: email=(candidate_email or "").strip().lower() if not email: raise HTTPException(status_code=422,detail="candidate_email is required") if not current_user: raise HTTPException(status_code=400,detail="created_by is required") original_name=(file_name or "").strip() or "resume.pdf" if file_bytes is not None: try: original_name=assert_pdf(original_name,content_type) except S3ServiceError as e: raise HTTPException(status_code=e.status_code,detail=e.message) from e if not file_bytes: raise HTTPException(status_code=422,detail="file is empty") parsed_linkedin=await parse_linkedin_url_from_cv(full_text) phone_fields=parse_phone( {"phone":(candidate_phone or "").strip()}, full_text or "", ) phone=phone_fields.get("phone") or "" data={ "candidate_email":email, "candidate_name":(candidate_name or "").strip(), "candidate_phone":phone, "job_post_id":job_post_id, "current_company":(current_company or "").strip(), "current_position":(current_position or "").strip(), "platform":(platform or "").strip(), "apply_via":"manual_upload", "experience":(experience or "").strip(), "status":(status or "").strip(), "referral_by":(referral_by or "").strip(), "file_name":original_name, # path filled after S3 succeeds; never leave a local orphan path here "file_path":(file_path or "").strip() if file_bytes is None else "", "full_text":full_text or "", "linkedin_url":parsed_linkedin, "created_by":current_user, } row=await Manual_UPLOAD_CANDIDATE.create_manual_upload_candidate(session=self.session,fields=data) if file_bytes is not None: try: uploaded=S3().upload_for_record( file_bytes, original_name, source=S3Source.MANUAL, record_id=row.id, owner_id=row.user_id, content_type=content_type, ) except S3ServiceError as e: await Manual_UPLOAD_CANDIDATE.delete_by_id(self.session,row.id) row=None raise HTTPException(status_code=e.status_code,detail=e.message) from e except Exception: await Manual_UPLOAD_CANDIDATE.delete_by_id(self.session,row.id) row=None raise row=await Manual_UPLOAD_CANDIDATE.set_file_path( self.session,row.id,uploaded["url"],file_name=uploaded.get("filename") or original_name, ) # Same permanent URL on candidates rows for this email+job (if scored already). try: await Candidates.sync_s3_file_path( self.session, email=row.candidate_email, job_id=row.job_post_id, file_path=row.file_path, ) except Exception: logger.exception("candidates.file_path sync failed for manual %s",row.id) await HistoryRecorder(self.session).record( HistoryEvent.CANDIDATE_CREATED.value, actor_id=current_user,user_id=row.user_id, manual_upload_candidate_id=row.id, entity_type="manual_upload_candidate",entity_id=row.id, to_value=row.candidate_email, description=(row.platform or "").strip() or "manual_upload",commit=True, ) if (row.file_name or "").strip() and (row.file_path or "").strip(): await HistoryRecorder(self.session).record( HistoryEvent.DOCUMENT_UPLOADED.value, actor_id=current_user,user_id=row.user_id, manual_upload_candidate_id=row.id, entity_type="document",entity_id=row.id, to_value=row.file_name,commit=True, ) return serialize_manual_upload_candidate(row) except HTTPException: raise except Exception as e: if row is not None: try: await Manual_UPLOAD_CANDIDATE.delete_by_id(self.session,row.id) except Exception: logger.exception("manual candidate rollback failed for %s",getattr(row,"id",None)) raise HTTPException(status_code=500,detail=str(e)) async def list_manager_candidates(self,current_user,limit=50,offset=0): """Candidates allocated to jobs this manager owns (requisition → job post).""" job_ids=await JobPosts.ids_for_manager(self.session,current_user.get("id")) if not job_ids: return [],0 inbox_rows=await Inbox.get_all(self.session,job_post_ids=job_ids) manual_rows=await Manual_UPLOAD_CANDIDATE.get_all(self.session,job_post_ids=job_ids) merged=[] seen=set() for row in inbox_rows: uid=row.get("user_id") payload=serialize_manager_candidate(row,source="inbox") if uid and uid not in seen: seen.add(uid) merged.append(payload) elif not uid: merged.append(payload) for row in manual_rows: uid=row.get("user_id") if uid and uid in seen: continue payload=serialize_manager_candidate(row,source="manual") if uid: seen.add(uid) merged.append(payload) merged.sort(key=lambda r: r.get("created_at") or "",reverse=True) total=len(merged) start=max(0,int(offset or 0)) cap=max(1,int(limit or 50)) return merged[start:start+cap],total async def get_candidate(self,user_id=None,limit=10,offset=0,search=None,current_user=None,assigned_job_post_id=None): try: if not user_id and is_hiring_manager(current_user): raise HTTPException(status_code=403,detail=MANAGER_SCOPE_DETAIL) if user_id: await assert_manager_candidate_access( self.session,current_user,user_id=user_id, ) detail=bool(user_id) list_job_ids=None if not detail: list_job_ids=await job_post_ids_for_candidate_list( self.session,current_user,assigned_job_post_id=assigned_job_post_id, ) if list_job_ids is not None and not list_job_ids: return [] rows=await Inbox.get_candidate_profile( session=self.session,user_id=user_id,limit=limit,offset=offset,search=search, job_post_ids=list_job_ids, ) if detail: records=rows if isinstance(rows,list) else ([rows] if rows else []) owned=await owned_job_ids_for_candidate_scope(self.session,current_user) if records and owned is not None: owned_set=set(owned) kept=[] for rec in records: msg=getattr(rec,"messages",None) jid=getattr(msg,"assigned_job_post_id",None) if msg is not None else None if jid and jid in owned_set: kept.append(rec) if kept: return await self.attach_profile_detail(kept) records=[] if records: return await self.attach_profile_detail(rows) # Manual uploads create users + manual_upload_candidate but no inbox # row — resolve the profile from that table instead of returning []. manual=await Manual_UPLOAD_CANDIDATE.get_by_user_id(self.session,user_id) if not manual: return [] if owned is not None: owned_set=set(owned) if not manual.job_post_id or manual.job_post_id not in owned_set: raise HTTPException(status_code=403,detail=_scope_detail(current_user)) user=await Users.get_user_by_id(self.session,user_id) job_post=None if manual.job_post_id: job_post=await JobPosts.get_job_post_by_id(self.session,str(manual.job_post_id)) payload=serialize_manual_candidate_profile(manual,user,job_post) from inbox.plugins import get_ats_score_for_manual_user score=await get_ats_score_for_manual_user(self.session,user_id,manual.job_post_id) if score: payload["ai_score"]=score["overall_score"] payload["recommendation"]=self._recommendation(score["overall_score"]) payload["scored_at"]=score["computed_at"] payload["candidate_id"]=score.get("candidate_id") if score.get("user_id") and not payload.get("user_id"): payload["user_id"]=score["user_id"] if score.get("job_post_id"): payload["scored_job_post_id"]=score["job_post_id"] return payload # List mode: inbox applications + manual/form applications (dedupe by user). inbox_payloads=await self.attach_job_posts(rows) if not isinstance(inbox_payloads,list): inbox_payloads=[inbox_payloads] if inbox_payloads else [] manual_rows=await Manual_UPLOAD_CANDIDATE.list_for_talent_pool( self.session,limit=limit,offset=0,search=search,job_post_ids=list_job_ids, ) seen={p.get("user_id") for p in inbox_payloads if p.get("user_id")} manual_payloads=[] for manual in manual_rows: uid=str(manual.user_id) if manual.user_id else None if uid and uid in seen: continue user=await Users.get_user_by_id(self.session,manual.user_id) if manual.user_id else None job_post=None if manual.job_post_id: job_post=await JobPosts.get_job_post_by_id(self.session,str(manual.job_post_id)) payload=serialize_manual_candidate_profile(manual,user,job_post) # List shape matches attach_job_posts: keep job_posts, drop heavy detail. manual_payloads.append({ "inbox_id":None, "manual_upload_candidate_id":payload["manual_upload_candidate_id"], "user_id":payload["user_id"], "candidate_id":None, "name":payload["name"], "email":payload["email"], "is_active":payload.get("is_active"), "message_id":None, "created_at":payload.get("created_at"), "application_status":payload.get("application_status"), "experience":payload.get("experience"), "current_employment":payload.get("current_employment"), "current_title":payload.get("current_title"), "resume_text":None, "suggested_job_post_ids":[], "assigned_job_post_id":payload.get("assigned_job_post_id"), "job_posts":payload.get("job_posts") or [], "assigned_job_post":payload.get("assigned_job_post"), "job_title":payload.get("job_title"), "recruiter":payload.get("recruiter"), "recruiter_id":payload.get("recruiter_id"), "source":payload.get("source"), "file_path":payload.get("file_path"), "ai_score":None, "recommendation":None, }) if uid: seen.add(uid) from inbox.plugins import get_ats_scores_for_users owners=[p.get("user_id") for p in manual_payloads if p.get("user_id")] ats=await get_ats_scores_for_users(self.session,owners) for payload in manual_payloads: row=ats.get(str(payload.get("user_id") or "")) if row and row.get("overall_score") is not None: payload["ai_score"]=row["overall_score"] payload["recommendation"]=row.get("band") or self._recommendation(row["overall_score"]) data=inbox_payloads+manual_payloads if assigned_job_post_id: job_id=str(assigned_job_post_id) data=[p for p in data if str(p.get("assigned_job_post_id") or "")==job_id] return data 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,current_user=None,assigned_job_post_id=None): try: job_post_ids=None if not user_id: job_post_ids=await job_post_ids_for_candidate_list( self.session,current_user,assigned_job_post_id=assigned_job_post_id, ) return await Inbox.count_candidate_profiles( session=self.session,user_id=user_id,search=search,job_post_ids=job_post_ids, ) except HTTPException: raise except Exception as e: raise HTTPException(status_code=500,detail=str(e)) async def update_candidate(self,user_id,payload,current_user=None): try: if not user_id: raise HTTPException(status_code=400,detail="user_id is required") await assert_manager_candidate_access(self.session,current_user,user_id=user_id) 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") first=records[0] old_favorite=getattr(first,"favorite",None) old_rating=getattr(first,"rating",None) for link in records: await Inbox.update_inbox(self.session,link.id,fields) if "favorite" in fields and fields["favorite"]!=old_favorite: await HistoryRecorder(self.session).record( HistoryEvent.FAVORITE_CHANGED.value, current_user=current_user,user_id=user_id,inbox_id=first.id, entity_type="candidate",entity_id=user_id, from_value=old_favorite,to_value=fields["favorite"],commit=True, ) if "rating" in fields and fields["rating"]!=old_rating: await HistoryRecorder(self.session).record( HistoryEvent.RATING_CHANGED.value, current_user=current_user,user_id=user_id,inbox_id=first.id, entity_type="candidate",entity_id=user_id, from_value=old_rating,to_value=fields["rating"],commit=True, ) 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)) @staticmethod def _attach_job_post(data,payload,*,as_assigned=False): """Merge one serialized job post onto a candidate payload. Pure, so the per-id path and the batched one below cannot drift. A missing post attaches nothing, matching the old early return. """ if not isinstance(data,dict) or not payload: return payload 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 async def _job_posts_by_id(self,ids): """Serialized job posts keyed by id — one query for a whole page of rows. active_only=False mirrors get_job_post_by_id, which filters neither flag: an application assigned to a closed post must keep its title. """ wanted=[] seen=set() for raw in ids or []: key=str(raw) if raw else None if not key or key in seen: continue seen.add(key) wanted.append(key) if not wanted: return {} rows=await JobPosts.get_by_ids(self.session,wanted,active_only=False) return {str(r.id):serialize_job_post(r) for r in rows} 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 return self._attach_job_post(data,serialize_job_post(job_post_data),as_assigned=as_assigned) 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. Three batched queries per page, not two per row — a 200-row pipeline board issued 600+ sequential job_posts round trips before. """ from inbox.plugins import get_ats_scores_for_users single=not isinstance(data,list) records=[data] if single else list(data or []) payloads=[] wanted=[] owners=[] for record in records: payload=serialize_candidate_profile(record) payload["job_posts"]=[] payload["assigned_job_post"]=None score,band=self._score_from_message(record) if score is not None: payload["ai_score"]=score payload["recommendation"]=band payloads.append(payload) if payload.get("user_id"): owners.append(payload["user_id"]) if payload.get("assigned_job_post_id"): wanted.append(payload["assigned_job_post_id"]) wanted.extend(payload.get("suggested_job_post_ids") or []) # ats_results is the source the inbox denorm above is copied FROM, so a # live score wins over it. A candidate with neither keeps ai_score None — # the list rows must not invent a number the scoring engine never produced. ats=await get_ats_scores_for_users(self.session,owners) for payload in payloads: row=ats.get(str(payload.get("user_id") or "")) if row and row.get("overall_score") is not None: payload["ai_score"]=row["overall_score"] payload["recommendation"]=row.get("band") or self._recommendation(row["overall_score"]) posts=await self._job_posts_by_id(wanted) # Copy per attach: two candidates on the same post held independent dicts # back when every row re-serialized its own. for payload in payloads: assigned_id=payload.get("assigned_job_post_id") if assigned_id: post=posts.get(str(assigned_id)) self._attach_job_post(payload,dict(post) if post else None,as_assigned=True) for job_id in payload.get("suggested_job_post_ids") or []: post=posts.get(str(job_id)) self._attach_job_post(payload,dict(post) if post else None) return payloads[0] if single else payloads 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: notes=[serialize_note(r) for r in await Notes.get_notes_by_user(self.session,uid)] # ATS score from inbox denorm / ats_results via Inbox.ats_id — never from # a Candidates join on message id. Keywords live on the scored Candidates # row: candidate_id when set, else email+job for the matched-user path. ats_ids=[r.ats_id for r in records if getattr(r,"ats_id",None)] ats_rows=await AtsResults.get_by_ids(self.session,ats_ids) if ats_ids else [] assigned_uid=AtsResults._as_uuid(base.get("assigned_job_post_id")) if base.get("assigned_job_post_id") else None chosen=None if assigned_uid is not None: chosen=next((a for a in ats_rows if a.job_post_id==assigned_uid),None) if chosen is None and ats_rows: chosen=max(ats_rows,key=lambda a:a.computed_at) if chosen is not None: base["ai_score"]=chosen.overall_score base["recommendation"]=chosen.band or self._recommendation(chosen.overall_score) base["scored_job_post_id"]=str(chosen.job_post_id) if chosen.job_post_id else None base["scored_at"]=chosen.computed_at.isoformat() if chosen.computed_at else None base["candidate_id"]=str(chosen.candidate_id) if chosen.candidate_id else None if chosen.user_id and not base.get("user_id"): base["user_id"]=str(chosen.user_id) scored=None if chosen.candidate_id: scored=await Candidates.get_candidate_by_id(self.session,str(chosen.candidate_id)) elif chosen.user_id: owner=await Users.get_user_by_id(self.session,str(chosen.user_id)) if owner is not None: scored=await Candidates.get_completed_by_email_job(self.session,owner.email,chosen.job_post_id) if scored is not None: base["matched_keywords"]=list(scored.matched_keywords or []) base["missing_keywords"]=list(scored.missing_keywords or []) base["summary_critique"]=scored.summary_critique else: for record in records: score,band=self._score_from_message(record) if score is not None: base["ai_score"]=score base["recommendation"]=band break 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 async def download_document(self,inbox_id=None,manual_upload_candidate_id=None,index=0): has_inbox=inbox_id is not None and str(inbox_id).strip()!="" has_manual=manual_upload_candidate_id is not None and str(manual_upload_candidate_id).strip()!="" if has_inbox==has_manual: raise HTTPException(status_code=404,detail="Not found") try: index=int(index or 0) except (TypeError,ValueError): raise HTTPException(status_code=404,detail="Not found") if index<0: raise HTTPException(status_code=404,detail="Not found") docs=[] if has_inbox: link=await Inbox.get_inbox_with_message(self.session,inbox_id) if not link or not link.messages: raise HTTPException(status_code=404,detail="Not found") docs=documents_from_message(link.messages.file_name,link.messages.file_path) else: row=await Manual_UPLOAD_CANDIDATE.get_by_id(self.session,manual_upload_candidate_id) if not row: raise HTTPException(status_code=404,detail="Not found") docs=documents_from_message(row.file_name,row.file_path) if index>=len(docs): raise HTTPException(status_code=404,detail="Not found") entry=docs[index] path=contained_download_path(entry.get("path")) if path is None: raise HTTPException(status_code=404,detail="Not found") name=(entry.get("name") or path.name).strip() or path.name return path,name async def list_matching(self,assigned=None,search=None,limit=10,offset=0): rows,total=await Manual_UPLOAD_CANDIDATE.list_matching( self.session,assigned=assigned,search=search,limit=limit,offset=offset, ) job_ids=[str(r.job_post_id) for r in rows if r.job_post_id] posts=await JobPosts.get_by_ids(self.session,job_ids,active_only=False) if job_ids else [] by_id={str(p.id):p for p in posts} data=[ serialize_matching_candidate(r,by_id.get(str(r.job_post_id)) if r.job_post_id else None) for r in rows ] return data,total async def get_matching(self,record_id): row=await Manual_UPLOAD_CANDIDATE.get_by_id(self.session,record_id) if not row or row.apply_via!="cv_bank": raise HTTPException(status_code=404,detail="CV not found") job_post=None if row.job_post_id: job_post=await JobPosts.get_job_post_by_id(self.session,str(row.job_post_id)) return serialize_matching_candidate(row,job_post) async def assign_matching(self,record_id,job_post_id,current_user=None): row=await Manual_UPLOAD_CANDIDATE.get_by_id(self.session,record_id) if not row or row.apply_via!="cv_bank": raise HTTPException(status_code=404,detail="CV not found") job_post=None if job_post_id not in (None,""): job_post=await JobPosts.get_job_post_by_id(self.session,str(job_post_id)) if not job_post or job_post.is_deleted: raise HTTPException(status_code=404,detail="Job post not found") was_unassigned=row.job_post_id is None row=await Manual_UPLOAD_CANDIDATE.assign_job_post(self.session,record_id,job_post_id) if not row: raise HTTPException(status_code=404,detail="CV not found") if job_post is None and row.job_post_id: job_post=await JobPosts.get_job_post_by_id(self.session,str(row.job_post_id)) if was_unassigned and row.job_post_id and row.user_id: title=(job_post.title if job_post else "") or str(row.job_post_id) await HistoryRecorder(self.session).record( HistoryEvent.CANDIDATE_CREATED.value, actor_id=current_user,user_id=row.user_id, manual_upload_candidate_id=row.id, entity_type="manual_upload_candidate",entity_id=row.id, to_value=title, description="cv_bank",commit=True, ) return serialize_matching_candidate(row,job_post)