diff --git a/backend/inbox/models.py b/backend/inbox/models.py index 0077623..f355f43 100644 --- a/backend/inbox/models.py +++ b/backend/inbox/models.py @@ -1,5 +1,6 @@ import logging import os +from shlex import join import uuid from datetime import datetime, timezone from typing import Any, List, Optional @@ -46,23 +47,14 @@ class Inbox(SQLModel, table=True): message_id: uuid.UUID | None = Field(default=None, foreign_key="inbox_messages.id") messages: Optional["Inbox_Messages"] = Relationship(back_populates="inbox") - # tz-AWARE, matching every other timestamp the analytics layer filters on. - # A naive column here made asyncpg reject the aware UTC bounds that - # analytics/views.py builds, so /analytics/hiring-trend and /analytics/kpis - # both 500'd before the query ever reached Postgres. created_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) updated_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) favorite: Optional[bool] = Field(default=False) rating: Optional[float] = Field(default=0.0) - # The CURRENT ats_results row for this application. Written together with the - # supersede chain in CandidateScoring._sync_inbox_ats: every completed inbox - # score inserts an ats_results row and repoints this at it, so the score is - # one direct id join away instead of a filter on ats_results.is_current. ats_id: uuid.UUID | None = Field(default=None, foreign_key="ats_results.id") - # selectin on one-to-many: joined would repeat the inbox row per child interviews: List["Interviews"] = Relationship( back_populates="inbox", sa_relationship_kwargs={"lazy": "selectin"}, @@ -80,6 +72,107 @@ class Inbox(SQLModel, table=True): sa_relationship_kwargs={"lazy": "joined"}, ) + @classmethod + async def get_all(cls,session:AsyncSession,job_post_id=None,limit=None,offset=0): + try: + from job.job_post.models import JobPosts + qry=( + select( + cls.id.label("inbox_id"), + cls.user_id, + Users.name, + Users.email, + Inbox_Messages.candidate_phone_number.label("phone"), + Inbox_Messages.assigned_job_post_id, + Inbox_Messages.application_status, + Inbox_Messages.current_employment, + Inbox_Messages.current_title, + Inbox_Messages.experience, + cls.created_at, + JobPosts.title, + AtsResults.id.label("ats_result_id"), + AtsResults.overall_score, + AtsResults.band, + AtsResults.job_post_id.label("ats_job_post_id"), + AtsResults.computed_at, + AtsResults.candidate_id, + AtsResults.user_id.label("ats_user_id"), + ) + .join(Users,cls.user_id==Users.id) + .join(Roles,Users.role_id==Roles.id) + .join(Inbox_Messages,cls.message_id==Inbox_Messages.id) + .join(JobPosts,Inbox_Messages.assigned_job_post_id==JobPosts.id) + .outerjoin(AtsResults,cls.ats_id==AtsResults.id) + .where(Inbox_Messages.assigned_job_post_id.is_not(None)) + .where(Roles.role_name==EnumRoles.CANDIDATE.value) + .order_by(cls.created_at.desc()) + ) + if job_post_id: + qry=qry.where(Inbox_Messages.assigned_job_post_id==job_post_id) + if limit is not None: + qry=qry.limit(limit).offset(offset) + result=await session.execute(qry) + rows=[] + for row in result.mappings().all(): + ats=None + if row["ats_result_id"] is not None: + ats={ + "id":str(row["ats_result_id"]), + "overall_score":row["overall_score"], + "band":row["band"] or None, + "job_post_id":str(row["ats_job_post_id"]) if row["ats_job_post_id"] else None, + "computed_at":row["computed_at"].isoformat() if row["computed_at"] else None, + "candidate_id":str(row["candidate_id"]) if row["candidate_id"] else None, + "user_id":str(row["ats_user_id"]) if row["ats_user_id"] else None, + } + status=row["application_status"] + rows.append({ + "inbox_id":row["inbox_id"], + "user_id":str(row["user_id"]) if row["user_id"] else None, + "name":row["name"], + "email":row["email"], + "application_status":status.value if status else None, + "phone":row["phone"], + "assigned_job_post_id":str(row["assigned_job_post_id"]) if row["assigned_job_post_id"] else None, + "title":row["title"] or None, + "current_employment":row["current_employment"] or None, + "current_title":row["current_title"] or None, + "experience":row["experience"] or None, + "created_at":row["created_at"].isoformat() if row["created_at"] else None, + "ats_result":ats, + }) + return rows + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + @classmethod + async def count_by_status(cls,session:AsyncSession,job_post_id=None): + try: + from job.job_post.models import JobPosts + qry=( + select(Inbox_Messages.application_status,func.count()) + .select_from(cls) + .join(Users,cls.user_id==Users.id) + .join(Roles,Users.role_id==Roles.id) + .join(Inbox_Messages,cls.message_id==Inbox_Messages.id) + .join(JobPosts,Inbox_Messages.assigned_job_post_id==JobPosts.id) + .where(Inbox_Messages.assigned_job_post_id.is_not(None)) + .where(Roles.role_name==EnumRoles.CANDIDATE.value) + .group_by(Inbox_Messages.application_status) + ) + if job_post_id: + qry=qry.where(Inbox_Messages.assigned_job_post_id==job_post_id) + result=await session.execute(qry) + counts={} + for status,n in result.all(): + key=status.value if hasattr(status,"value") else (str(status) if status else None) + if not key: + continue + counts[key]=int(n or 0) + return counts + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + @classmethod def _candidate_search_filter(cls, search: str): pattern = f"%{search}%" @@ -507,6 +600,18 @@ class Inbox_Messages(SQLModel, table=True): result = await session.execute(select(cls).where(cls.id == uid)) return result.scalars().first() + @classmethod + async def set_ats_score(cls, session: AsyncSession, record_id, score, band): + row = await cls.get_inbox_message_by_id(session, record_id) + if not row: + return None + row.ats_score = float(score) + row.ats_band = band or "" + session.add(row) + await session.commit() + await session.refresh(row) + return row + @classmethod async def set_assigned_job_post(cls, session: AsyncSession, record_id, job_post_id): """Set or clear assigned_job_post_id; returns the row or None if missing.""" @@ -608,9 +713,11 @@ class AtsResults(SQLModel, table=True): id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True) # Inbox applications link here; NULL for upload-sourced scores, which have no - # inbox row. candidate_id links every score to the candidates row it scored. + # inbox row. Identity is XOR: matching users.email -> user_id (candidate_id + # NULL); otherwise candidate_id (user_id NULL). Blank email is the latter. inbox_id: int | None = Field(default=None, index=True, foreign_key="inbox.id") candidate_id: uuid.UUID | None = Field(default=None, index=True, foreign_key="candidates.id") + user_id: uuid.UUID | None = Field(default=None, index=True, foreign_key="users.id") job_post_id: uuid.UUID | None = Field(default=None, foreign_key="job_posts.id") overall_score: float = Field(default=0.0) band: str = Field(default="") @@ -641,6 +748,41 @@ class AtsResults(SQLModel, table=True): ) return result.scalars().first() + @classmethod + async def get_for_inbox_job(cls, session: AsyncSession, inbox_id, job_post_id): + """Any score for this application against this job — current or superseded. + + Inbox-side idempotency is (inbox_id, job_post_id) only. candidate_id may + be NULL when the CV email matched a user; do not key off it here. + """ + jid = cls._as_uuid(job_post_id) + if jid is None: + return None + result = await session.execute( + select(cls) + .where(cls.inbox_id == int(inbox_id), cls.job_post_id == jid) + .order_by(cls.computed_at.desc()) + ) + return result.scalars().first() + + @classmethod + async def resolve_identity(cls, session: AsyncSession, email, candidate_id): + """XOR identity for a score row from the scored candidate's email. + + Matching users.email (case-insensitive) -> user_id, candidate_id NULL. + Missing/blank email or no user -> candidate_id, user_id NULL. + """ + normalized = (email or "").strip().lower() + if normalized: + user_id = ( + await session.execute( + select(Users.id).where(func.lower(Users.email) == normalized) + ) + ).scalar_one_or_none() + if user_id is not None: + return {"candidate_id": None, "user_id": user_id} + return {"candidate_id": candidate_id, "user_id": None} + @classmethod async def get_current_for_candidate(cls, session: AsyncSession, candidate_id): """Current row for an upload-sourced score, chained per candidates row — @@ -653,9 +795,52 @@ class AtsResults(SQLModel, table=True): ) return result.scalars().first() + @classmethod + async def get_current_for_user(cls, session: AsyncSession, user_id, job_post_id=None): + """Current upload-sourced score for a matched user, scoped per job.""" + uid = cls._as_uuid(user_id) + if uid is None: + return None + qry = select(cls).where(cls.user_id == uid, cls.is_current == True) # noqa: E712 + jid = cls._as_uuid(job_post_id) if job_post_id is not None else None + if jid is not None: + qry = qry.where(cls.job_post_id == jid) + result = await session.execute(qry.order_by(cls.computed_at.desc())) + return result.scalars().first() + @classmethod async def insert_result(cls, session: AsyncSession, fields: dict): + """Insert a score row and supersede the previous current one. + + Inbox scores chain on inbox_id and repoint inbox.ats_id (candidate_id + may be NULL). Upload scores chain on candidate_id, or on (user_id, + job_post_id) when identity resolved to a user. Flush the INSERT first: + with no relationship() edge the unit of work emits the UPDATEs first, + and the FKs reject a pointer to a row not yet inserted.""" + inbox_id = fields.get("inbox_id") + candidate_id = fields.get("candidate_id") + user_id = fields.get("user_id") + job_post_id = fields.get("job_post_id") + if inbox_id is not None: + prev = await cls.get_current_for_inbox(session, inbox_id) + elif candidate_id is not None: + prev = await cls.get_current_for_candidate(session, candidate_id) + elif user_id is not None: + prev = await cls.get_current_for_user(session, user_id, job_post_id) + else: + prev = None row = cls(**fields) session.add(row) + await session.flush() + if prev is not None: + prev.is_current = False + prev.superseded_by_id = row.id + session.add(prev) + if inbox_id is not None: + link = await Inbox.get_inbox_by_id(session, inbox_id) + if link is not None: + link.ats_id = row.id + link.updated_at = _now() + session.add(link) await session.commit() return row diff --git a/backend/inbox/plugins.py b/backend/inbox/plugins.py index 7af0ec8..45a2006 100644 --- a/backend/inbox/plugins.py +++ b/backend/inbox/plugins.py @@ -5,13 +5,17 @@ from __future__ import annotations import base64 import os import re +import uuid from pathlib import Path from urllib.parse import quote import httpx from dotenv import load_dotenv +from sqlalchemy.ext.asyncio import AsyncSession +from sqlmodel import select -from inbox.models import Inbox_Messages +from inbox.models import AtsResults, Inbox, Inbox_Messages +from job.candidate.models import Candidates, Manual_UPLOAD_CANDIDATE from job.candidate.views import FileRead load_dotenv() @@ -161,3 +165,128 @@ async def extract_resume_text(file_paths:list[str]) -> tuple[str,str]: if not texts: return "","; ".join(errors) if errors else "no text extracted from PDF" return "\n\n---\n\n".join(texts),"" + + +def _ats_score_payload(row): + if row is None: + return None + return { + "overall_score":row.overall_score, + "band":row.band or None, + "job_post_id":str(row.job_post_id) if row.job_post_id else None, + "computed_at":row.computed_at.isoformat() if row.computed_at else None, + "candidate_id":str(row.candidate_id) if row.candidate_id else None, + "user_id":str(row.user_id) if row.user_id else None, + } + + +async def get_ats_score_for_user(session:AsyncSession,user_id,job_post_id=None): + """Current ats_results overall score for a candidate -> dict, or None. + + Prefer ats_results.user_id (CV email matched that user). Fall back to the + inbox.user_id join for scores whose email did not match any user, where + candidate_id is set and user_id is NULL. + + Pass job_post_id to pin one application when a candidate has several; + without it the newest current score across their applications wins. + """ + try: + uid=uuid.UUID(str(user_id)) + except (TypeError,ValueError): + return None + + direct=( + select(AtsResults) + .where(AtsResults.user_id==uid,AtsResults.is_current==True) # noqa: E712 + .order_by(AtsResults.computed_at.desc()) + ) + if job_post_id: + try: + direct=direct.where(AtsResults.job_post_id==uuid.UUID(str(job_post_id))) + except (TypeError,ValueError): + return None + row=(await session.execute(direct)).scalars().first() + if row is not None: + return _ats_score_payload(row) + + qry=( + select(AtsResults) + .join(Inbox,AtsResults.inbox_id==Inbox.id) + .join(Inbox_Messages,Inbox.message_id==Inbox_Messages.id) + .where( + Inbox.user_id==uid, + Inbox_Messages.assigned_job_post_id.is_not(None), + AtsResults.job_post_id==Inbox_Messages.assigned_job_post_id, + AtsResults.is_current==True, # noqa: E712 + ) + .order_by(AtsResults.computed_at.desc()) + ) + if job_post_id: + try: + qry=qry.where(Inbox_Messages.assigned_job_post_id==uuid.UUID(str(job_post_id))) + except (TypeError,ValueError): + return None + + return _ats_score_payload((await session.execute(qry)).scalars().first()) + + +async def get_ats_score_for_manual_user(session:AsyncSession,user_id,job_post_id=None): + """Current ats_results overall score for an Add Candidate user -> dict, or None. + + Prefer ats_results.user_id (Add Candidate always creates a users row, so a + later score against that email lands on user_id). Fall back to the + email+candidate_id join for rows written before user_id existed. + apply_via=manual_upload is the Add Candidate gate; /import never writes + that table. + + Pass job_post_id to pin one application when a candidate has several; + without it the newest current score across their applications wins. + """ + try: + uid=uuid.UUID(str(user_id)) + except (TypeError,ValueError): + return None + + direct=( + select(AtsResults) + .where(AtsResults.user_id==uid,AtsResults.is_current==True) # noqa: E712 + .order_by(AtsResults.computed_at.desc()) + ) + if job_post_id: + try: + jid=uuid.UUID(str(job_post_id)) + except (TypeError,ValueError): + return None + direct=direct.where(AtsResults.job_post_id==jid) + row=(await session.execute(direct)).scalars().first() + if row is not None: + return _ats_score_payload(row) + + qry=( + select(AtsResults) + .join(Candidates,AtsResults.candidate_id==Candidates.id) + .join( + Manual_UPLOAD_CANDIDATE, + (Candidates.job_id==Manual_UPLOAD_CANDIDATE.job_post_id) + &(Candidates.candidate_email==Manual_UPLOAD_CANDIDATE.candidate_email), + ) + .where( + Manual_UPLOAD_CANDIDATE.user_id==uid, + Manual_UPLOAD_CANDIDATE.apply_via=="manual_upload", + Candidates.status=="completed", + AtsResults.job_post_id==Manual_UPLOAD_CANDIDATE.job_post_id, + AtsResults.is_current==True, # noqa: E712 + ) + .order_by(AtsResults.computed_at.desc()) + ) + if job_post_id: + try: + jid=uuid.UUID(str(job_post_id)) + except (TypeError,ValueError): + return None + qry=qry.where( + Manual_UPLOAD_CANDIDATE.job_post_id==jid, + AtsResults.job_post_id==jid, + ) + + return _ats_score_payload((await session.execute(qry)).scalars().first()) diff --git a/backend/inbox/tasks.py b/backend/inbox/tasks.py index 0ce2f7e..ef25a70 100644 --- a/backend/inbox/tasks.py +++ b/backend/inbox/tasks.py @@ -3,15 +3,15 @@ from __future__ import annotations import logging +import uuid from datetime import datetime,timezone from fastapi import HTTPException -from sqlalchemy import select from agent.execute_agent import run_agent from db_setup import session_scope from employment_agent.execute_agent import run_employment_agent -from inbox.models import Inbox_Messages +from inbox.models import Inbox_Messages,Inbox,AtsResults from inbox.plugins import extract_phone,extract_resume_text from job.job_post.models import JobPosts from job.job_post.serializers import serialize_job_post @@ -26,29 +26,27 @@ async def score_message_against_job(record_id:str,job_id:str) -> dict: """ATS-score one inbox CV against one job post — the no-upload path. The decoded attachment already on disk is the CV; the job post in the - database is the JD. Idempotent: a (message, job) pair with a completed - score is never paid for twice; re-runs are a no-op. + database is the JD. Idempotent: a (message, job) pair with an ats_results + row is never paid for twice; re-runs are a no-op. """ # Lazy imports: inbox.plugins imports job.candidate.views, so a top-level # import here would be circular. - from job.candidate.models import Candidates from job.candidate.views import CandidateScoring - mid=Candidates._as_uuid(record_id) - jid=Candidates._as_uuid(job_id) - if mid is None or jid is None: + try: + mid=uuid.UUID(str(record_id)) + jid=uuid.UUID(str(job_id)) + except ValueError: raise PermanentTaskError("record_id and job_id must be uuids") async with session_scope() as session: - existing=await session.execute( - select(Candidates).where( - Candidates.inbox_message_id==mid, - Candidates.job_id==jid, - Candidates.status=="completed", - ) - ) - if existing.scalars().first() is not None: - return {"status":"already_scored"} + link=await Inbox.get_inbox_by_message_id(session,mid) + if link is not None: + # (inbox, job) only — candidate_id is NULL when the CV email matched + # a user, so already_scored must not depend on it. + existing=await AtsResults.get_for_inbox_job(session,link.id,jid) + if existing is not None: + return {"status":"already_scored"} job=await JobPosts.get_job_post_by_id(session,job_id) if job is None or job.is_deleted: raise PermanentTaskError("job post missing or deleted") diff --git a/backend/job/app.py b/backend/job/app.py index cfc8b0e..acd57c8 100644 --- a/backend/job/app.py +++ b/backend/job/app.py @@ -21,7 +21,8 @@ from dotenv import load_dotenv from datetime import datetime, time, timezone from pydantic import BaseModel from uuid import UUID - +from typing import Optional +import uuid load_dotenv() logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) @@ -89,7 +90,8 @@ class FeedbackUpdate(BaseModel): class StageChange(BaseModel): - inbox_id: int + inbox_id: int | None = None + manual_upload_id: UUID | None = None to_stage: str change_reason: str | None = None @@ -666,7 +668,8 @@ async def change_candidate_stage( try: service=Pipeline(session=session) data=await service.change_stage( - payload.inbox_id,payload.to_stage,current_user,change_reason=payload.change_reason, + payload.to_stage,current_user,inbox_id=payload.inbox_id, + manual_upload_id=payload.manual_upload_id,change_reason=payload.change_reason, ) return JSONResponse(content={"data":data,"total":1,"status_code":200}) except HTTPException: @@ -674,17 +677,52 @@ async def change_candidate_stage( except Exception as e: raise HTTPException(status_code=500,detail=str(e)) +@router.get("/pipeline/candidates/fetch") +async def fetch_pipeline_candidates( + job_post_id:Optional[uuid.UUID]=Query(None), + limit:int=Query(200,ge=1,le=1000), + offset:int=Query(0,ge=0), + current_user: dict = Depends(require_permission(PermissionTag.PIPELINE_VIEW)), + session: AsyncSession = Depends(get_session), +): + try: + service=Pipeline(session=session) + result=await service.get_all(job_post_id=job_post_id,limit=limit,offset=offset) + return JSONResponse(content={**result,"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + +@router.get("/pipeline/candidate/score/fetch") +async def fetch_pipeline_candidate_score( + user_id:uuid.UUID=Query(...), + job_post_id:uuid.UUID=Query(...), + current_user: dict = Depends(require_permission(PermissionTag.PIPELINE_VIEW)), + session: AsyncSession = Depends(get_session), +): + try: + service=Pipeline(session=session) + data=await service.get_pipeline_candidates(user_id=user_id,job_post_id=job_post_id) + return JSONResponse(content={"data":data,"total":1,"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + @router.get("/pipeline/transitions/fetch") async def fetch_pipeline_transitions( transition_id:str=Query(None), inbox_id:int=Query(None), + manual_upload_id:UUID=Query(None), current_user: dict = Depends(require_permission(PermissionTag.PIPELINE_VIEW)), session: AsyncSession = Depends(get_session), ): try: service=Pipeline(session=session) - data=await service.get_transitions(inbox_id=inbox_id,transition_id=transition_id) + data=await service.get_transitions(inbox_id=inbox_id,transition_id=transition_id,manual_upload_id=manual_upload_id) total=1 if isinstance(data,dict) else len(data) return JSONResponse(content={"data":data,"total":total,"status_code":200}) except HTTPException: diff --git a/backend/job/candidate/models.py b/backend/job/candidate/models.py index 72b84bc..4119c54 100644 --- a/backend/job/candidate/models.py +++ b/backend/job/candidate/models.py @@ -2,6 +2,7 @@ import uuid from datetime import datetime, timezone from typing import TYPE_CHECKING, List, Optional +from fastapi import HTTPException from sqlalchemy import JSON, DateTime, func, UniqueConstraint from sqlalchemy.exc import IntegrityError from sqlalchemy.ext.asyncio import AsyncSession @@ -35,33 +36,122 @@ class Manual_UPLOAD_CANDIDATE(SQLModel, table=True): # from job_posts.title — that is the role they applied to, not their own. # Same ALTER-on-a-populated-table reasoning as referral_by below. current_position: str = Field(default="", sa_column_kwargs={"server_default": ""}) + # Add Candidate only (POST /candidate/create/candidate). /import never writes this table. + apply_via:Optional[str]=Field(nullable=True) user_id: uuid.UUID | None = Field(default=None, foreign_key="users.id") platform: str = Field(default="") created_by: uuid.UUID | None = Field(default=None, foreign_key="users.id") experience: str = Field(default="") + # Candidate_application_Status value (PENDING, SCREENING, …). Empty reads as Applied. status: str = Field(default="") # Free text, not a users FK: a referrer is often someone outside the system - # (a client, a former colleague), and recruiters type whatever the candidate - # told them. "" rather than NULL keeps it consistent with the columns above. - # - # server_default is load-bearing and NOT decoration, unlike the columns above - # — they arrived with the CREATE TABLE, this one arrives as an ALTER. The - # startup autogenerate would emit `ADD COLUMN referral_by VARCHAR NOT NULL`, - # which Postgres rejects outright on a table that already holds rows. The - # DEFAULT backfills them. Pass the bare "" — SQLAlchemy quotes a plain string - # into DEFAULT '', whereas "''" would render DEFAULT '''''' instead. referral_by: str = Field(default="", sa_column_kwargs={"server_default": ""}) - # The CV as uploaded: file_name is the recruiter-facing original, file_path - # the absolute location under inbox/decoded_attachments. They differ on - # purpose — the stored basename is uniquified so two candidates uploading - # "resume.pdf" cannot overwrite one another (see FileRead.save_manual_upload). - # Same ALTER-on-a-populated-table reasoning as referral_by above, so both - # carry a server default. file_name: str = Field(default="", sa_column_kwargs={"server_default": ""}) file_path: str = Field(default="", sa_column_kwargs={"server_default": ""}) created_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) updated_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) + @classmethod + async def get_all(cls, session: AsyncSession, job_post_id=None, limit=None, offset=0): + try: + from inbox.models import AtsResults + from users.models import Users + from job.job_post.models import JobPosts + qry=( + select( + cls.id, + cls.candidate_email, + cls.user_id, + Users.email, + cls.job_post_id, + Users.name, + cls.candidate_phone, + JobPosts.title, + cls.status, + cls.current_company, + cls.current_position, + cls.experience, + cls.created_at, + cls.updated_at, + AtsResults.id.label("ats_result_id"), + AtsResults.overall_score, + AtsResults.band, + AtsResults.job_post_id.label("ats_job_post_id"), + AtsResults.computed_at, + AtsResults.candidate_id, + AtsResults.user_id.label("ats_user_id"), + ) + .join(Users,cls.user_id==Users.id) + .join(JobPosts,cls.job_post_id==JobPosts.id) + .outerjoin( + AtsResults, + (Users.id==AtsResults.user_id) + &(AtsResults.job_post_id==cls.job_post_id) + &(AtsResults.is_current==True), # noqa: E712 + ) + .order_by(cls.created_at.desc()) + ) + if job_post_id: + qry=qry.where(cls.job_post_id==job_post_id) + if limit is not None: + qry=qry.limit(limit).offset(offset) + result=await session.execute(qry) + rows=[] + for row in result.mappings().all(): + ats=None + if row["ats_result_id"] is not None: + ats={ + "id":str(row["ats_result_id"]), + "overall_score":row["overall_score"], + "band":row["band"] or None, + "job_post_id":str(row["ats_job_post_id"]) if row["ats_job_post_id"] else None, + "computed_at":row["computed_at"].isoformat() if row["computed_at"] else None, + "candidate_id":str(row["candidate_id"]) if row["candidate_id"] else None, + "user_id":str(row["ats_user_id"]) if row["ats_user_id"] else None, + } + rows.append({ + "id":str(row["id"]) if row["id"] else None, + "candidate_email":row["candidate_email"], + "user_id":str(row["user_id"]) if row["user_id"] else None, + "email":row["email"], + "job_post_id":str(row["job_post_id"]) if row["job_post_id"] else None, + "name":row["name"], + "candidate_phone":row["candidate_phone"], + "title":row["title"] or None, + "application_status":row["status"] or None, + "current_company":row["current_company"] or None, + "current_position":row["current_position"] or None, + "experience":row["experience"] or None, + "created_at":row["created_at"].isoformat() if row["created_at"] else None, + "updated_at":row["updated_at"].isoformat() if row["updated_at"] else None, + "ats_result":ats, + }) + return rows + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + @classmethod + async def count_by_status(cls, session: AsyncSession, job_post_id=None): + try: + from users.models import Users + from job.job_post.models import JobPosts + qry=( + select(cls.status,func.count()) + .select_from(cls) + .join(Users,cls.user_id==Users.id) + .join(JobPosts,cls.job_post_id==JobPosts.id) + .group_by(cls.status) + ) + if job_post_id: + qry=qry.where(cls.job_post_id==job_post_id) + result=await session.execute(qry) + counts={} + for status,n in result.all(): + key=(status or "").strip() or "UNKNOWN" + counts[key]=counts.get(key,0)+int(n or 0) + return counts + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) @staticmethod def _as_uuid(record_id) -> uuid.UUID | None: if record_id in (None, ""): @@ -103,11 +193,12 @@ class Manual_UPLOAD_CANDIDATE(SQLModel, table=True): full_text=fields.get("full_text") or "", current_company=(fields.get("current_company") or "").strip(), current_position=(fields.get("current_position") or "").strip(), + apply_via="manual_upload", user_id=user.id, platform=(fields.get("platform") or "").strip(), created_by=cls._as_uuid(fields.get("created_by")), experience=(fields.get("experience") or "").strip(), - status=(fields.get("status") or "").strip(), + status=(fields.get("status") or "").strip() or "PENDING", referral_by=(fields.get("referral_by") or "").strip(), file_name=(fields.get("file_name") or "").strip(), file_path=(fields.get("file_path") or "").strip(), @@ -127,29 +218,28 @@ class Manual_UPLOAD_CANDIDATE(SQLModel, table=True): ) return result.scalars().first() + @classmethod + async def get_by_id(cls, session: AsyncSession, record_id): + uid = cls._as_uuid(record_id) + if uid is None: + return None + result = await session.execute(select(cls).where(cls.id == uid)) + return result.scalars().first() + class Candidates(SQLModel, table=True): - """One scored (or failed-to-score) CV against one job post. - The dedupe key is (job_id, content_sha256), not the filename: inbox attachments - are stored by basename so different candidates can collide on "resume.pdf", while - identical bytes can arrive via both upload and email. Re-scoring the same bytes - against the same job updates the existing row (fresh model output, updated_at - bumped) instead of duplicating it. content_sha256 is NULL when the file bytes - were never readable (missing on disk); NULLs never conflict in the unique index. - """ __tablename__ = "candidates" __table_args__ = (UniqueConstraint("job_id", "content_sha256"),) id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True) job_id: uuid.UUID = Field(foreign_key="job_posts.id", index=True) - source: str = Field(default="upload") # "upload" | "inbox" - inbox_message_id: uuid.UUID | None = Field(default=None, foreign_key="inbox_messages.id") + source: str = Field(default="upload") # "upload" | "inbox" origin metadata, not an FK filename: str file_path: str | None = Field(default=None) # decoded-attachment path (inbox only) content_sha256: str | None = Field(default=None, index=True) - + candidate_email: str | None = Field(default=None) candidate_name: str | None = Field(default=None) job_title: str | None = Field(default=None) current_company: str | None = Field(default=None) @@ -204,6 +294,25 @@ class Candidates(SQLModel, table=True): result = await session.execute(statement) return result.scalars().all() + @classmethod + async def get_completed_by_email_job(cls, session: AsyncSession, email, job_id): + """Newest completed score for this email against this job — keywords when + ats_results.candidate_id is NULL (matched-user identity).""" + normalized = (email or "").strip().lower() + jid = cls._as_uuid(job_id) + if not normalized or jid is None: + return None + result = await session.execute( + select(cls) + .where( + func.lower(cls.candidate_email) == normalized, + cls.job_id == jid, + cls.status == "completed", + ) + .order_by(cls.updated_at.desc()) + ) + return result.scalars().first() + @classmethod async def upsert_candidate(cls, session: AsyncSession, fields: dict): existing = None @@ -525,16 +634,20 @@ class Feedback(SQLModel, table=True): class ApplicationStageTransitions(SQLModel, table=True): - """Temporal history of inbox_messages.application_status changes. + """Temporal history of application stage changes. - valid_from / valid_to make time-in-stage a subtraction rather than a window - function. NULL valid_to means the stage is still current. + Inbox moves write inbox_messages.application_status; manual-upload moves + write manual_upload_candidate.status. Exactly one of inbox_id / + manual_upload_candidate_id is set. NULL valid_to means the stage is current. """ __tablename__ = "application_stage_transitions" id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True) - inbox_id: int = Field(index=True, foreign_key="inbox.id") + inbox_id: int | None = Field(default=None, index=True, foreign_key="inbox.id") + manual_upload_candidate_id: uuid.UUID | None = Field( + default=None, index=True, foreign_key="manual_upload_candidate.id" + ) from_stage: str | None = Field(default=None) to_stage: str valid_from: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) @@ -569,12 +682,28 @@ class ApplicationStageTransitions(SQLModel, table=True): return list(result.scalars().all()) @classmethod - async def get_open_transition(cls, session: AsyncSession, inbox_id: int): + async def fetch_by_manual(cls, session: AsyncSession, record_id): + uid = cls._as_uuid(record_id) + if uid is None: + return [] result = await session.execute( - select(cls) - .where(cls.inbox_id == int(inbox_id), cls.valid_to.is_(None)) - .order_by(cls.valid_from.desc()) + select(cls).where(cls.manual_upload_candidate_id == uid).order_by(cls.valid_from.desc()) ) + return list(result.scalars().all()) + + @classmethod + async def get_open_transition(cls, session: AsyncSession, inbox_id=None, manual_upload_candidate_id=None): + statement = select(cls).where(cls.valid_to.is_(None)).order_by(cls.valid_from.desc()) + if inbox_id is not None: + statement = statement.where(cls.inbox_id == int(inbox_id)) + elif manual_upload_candidate_id is not None: + uid = cls._as_uuid(manual_upload_candidate_id) + if uid is None: + return None + statement = statement.where(cls.manual_upload_candidate_id == uid) + else: + return None + result = await session.execute(statement) return result.scalars().first() @classmethod @@ -587,8 +716,10 @@ class ApplicationStageTransitions(SQLModel, table=True): return row @classmethod - async def close_open(cls, session: AsyncSession, inbox_id: int, *, at: datetime | None = None, commit: bool = False): - row = await cls.get_open_transition(session, inbox_id) + async def close_open(cls, session: AsyncSession, inbox_id=None, *, manual_upload_candidate_id=None, at: datetime | None = None, commit: bool = False): + row = await cls.get_open_transition( + session, inbox_id=inbox_id, manual_upload_candidate_id=manual_upload_candidate_id + ) if not row: return None row.valid_to = at or _now() diff --git a/backend/job/candidate/plugins.py b/backend/job/candidate/plugins.py index 2b68fac..cb433a6 100644 --- a/backend/job/candidate/plugins.py +++ b/backend/job/candidate/plugins.py @@ -94,6 +94,49 @@ def build_job_description(job) -> str: return "\n".join(lines) +def candidate_base_fields(source): + return { + "filename": source["safe_name"], + "file_path": source["file_path"], + "content_sha256": source["sha256"], + "candidate_email": source.get("candidate_email"), + } + + +def candidate_failed_fields(source, code, message): + return { + **candidate_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, + } + + +def candidate_completed_fields(source, result): + return { + **candidate_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, + } + + @normalize_unicode @despace_line def normalize_spaced_text(text) -> str: diff --git a/backend/job/candidate/serializers.py b/backend/job/candidate/serializers.py index f5b2e07..7ba7351 100644 --- a/backend/job/candidate/serializers.py +++ b/backend/job/candidate/serializers.py @@ -11,11 +11,11 @@ def serialize_candidate(row) -> dict: return { "id": str(row.id), "job_id": str(row.job_id), - "inbox_message_id": str(row.inbox_message_id) if row.inbox_message_id else None, "source": row.source, "filename": row.filename, "file_path": row.file_path, "content_sha256": row.content_sha256, + "candidate_email": row.candidate_email, "candidate_name": row.candidate_name, "job_title": row.job_title, "current_company": row.current_company, @@ -44,6 +44,7 @@ def serialize_manual_upload_candidate(row) -> Dict[str,Any]: "full_text":row.full_text, "current_company":row.current_company, "current_position":row.current_position, + "apply_via":row.apply_via, "user_id":str(row.user_id) if row.user_id else None, "platform":row.platform, "created_by":str(row.created_by) if row.created_by else None, @@ -72,6 +73,7 @@ def serialize_candidate_profile( payload = { "inbox_id": link.id, "user_id": str(link.user_id) if link.user_id else None, + "candidate_id": None, "name": user.name if user else None, "email": user.email if user else None, "is_active": user.is_active if user else None, @@ -138,6 +140,7 @@ def serialize_manual_candidate_profile(row, user, job_post) -> Dict[str, Any]: return { "inbox_id": None, "user_id": str(user.id) if user else (str(row.user_id) if row.user_id else None), + "candidate_id": None, "name": (user.name if user else None) or row.candidate_name or None, "email": (user.email if user else None) or row.candidate_email or None, "is_active": user.is_active if user else None, diff --git a/backend/job/candidate/views.py b/backend/job/candidate/views.py index e2bfb33..0ea18f8 100644 --- a/backend/job/candidate/views.py +++ b/backend/job/candidate/views.py @@ -17,6 +17,8 @@ from job.candidate.models import Candidates from job.candidate.plugins import ( FILE_NOT_FOUND, build_job_description, + candidate_completed_fields, + candidate_failed_fields, get_scorer, get_scoring_settings, normalize_spaced_text, @@ -272,7 +274,6 @@ class CandidateScoring: "filename":filename or "resume.pdf", "data":data, "file_path":None, - "inbox_message_id":None, "precheck":None, } if not (filename or "").lower().endswith(".pdf"): @@ -301,7 +302,8 @@ class CandidateScoring: "filename":path.name, "data":None, "file_path":str(path), - "inbox_message_id":row.id, + "inbox_message_id":row.id, # call-scoped; not persisted on Candidates + "candidate_email":(row.message_from or "").strip().lower() or None, "precheck":None, } suffix=path.suffix.lower() @@ -344,10 +346,24 @@ class CandidateScoring: 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) + 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)): + rows.append(await Candidates.upsert_candidate(self.session,{**fields_by_slot[slot],**common})) + await self._sync_ats_results(source_kind,job,rows,sources) + 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] - # Slot-indexed like app/api/routes.py: results merge back by position, never - # by filename — inbox attachments can share a basename. - results_by_slot={} + 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"]) @@ -355,21 +371,20 @@ class CandidateScoring: 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) + fields_by_slot[slot]=candidate_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)) + if not source.get("candidate_email"): + detected,_=extract_candidate_email(resume.text) + source["candidate_email"]=detected except ATSError as exc: - results_by_slot[slot]=self._failed_fields(source,exc.error_code,exc.public_message) + 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, @@ -379,178 +394,80 @@ class CandidateScoring: 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, - } + fields_by_slot[slot]=candidate_completed_fields(source,result) else: - results_by_slot[slot]=self._failed_fields(source,result.error_code,result.error_message) + fields_by_slot[slot]=candidate_failed_fields(source,result.error_code,result.error_message) + return fields_by_slot - 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)) - - # Every completed score lands in ats_results. Inbox scores additionally - # denormalise onto inbox_messages / inbox (one sync per message: a - # multi-attachment mail keeps its best completed score); upload scores - # chain per candidates row instead — they have no inbox application. + async def _sync_ats_results(self,source_kind,job,rows,sources): if source_kind=="inbox": best={} - for row in rows: - if row.status=="completed" and row.inbox_message_id: - cur=best.get(row.inbox_message_id) + 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[row.inbox_message_id]=row + best[mid]=row for message_id,row in best.items(): try: await self._sync_inbox_ats(message_id,job,row) except Exception: - # The candidates row is the primary outcome and is already - # committed; a denorm failure must not fail the scoring call. await self.session.rollback() logger.exception("inbox ATS denorm failed for message %s",message_id) - else: - for row in rows: - if row.status!="completed": - continue - try: - await self._sync_upload_ats(job,row) - except Exception: - await self.session.rollback() - logger.exception("upload ATS history failed for candidate %s",row.id) - - # 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, - } + return + for row in rows: + if row.status!="completed": + continue + try: + await self._sync_upload_ats(job,row) + 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): - """Land a completed score on the inbox tables (README "What is missing"). + """Land a completed score on inbox_messages / inbox / ats_results. - inbox_messages.ats_score / ats_band are what serialize_application renders - on the Applications tab; ats_results keeps the per-application history that - candidates' in-place upsert cannot. Precedence mirrors _recommendation: a - score against the assigned job always wins the denormalised columns, any - other job's score only lands while no completed assigned-job score exists. - The history row is appended regardless — it records the scoring event. + 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 - if assigned and str(assigned)!=str(job.id): - outranked=await self.session.execute( - select(Candidates).where( - Candidates.inbox_message_id==msg.id, - Candidates.job_id==assigned, - Candidates.status=="completed", - ) - ) - denorm=outranked.scalars().first() is None - if denorm: - msg.ats_score=float(row.match_score) - msg.ats_band=band - self.session.add(msg) - - # ats_results hangs off the inbox JOIN row (int PK), which only exists once - # the sender is linked to a users account; without it there is no history row. link=await Inbox.get_inbox_by_message_id(self.session,message_id) - if link is not None: - prev=await AtsResults.get_current_for_inbox(self.session,link.id) - entry=AtsResults( - inbox_id=link.id, - candidate_id=row.id, - job_post_id=job.id, - overall_score=float(row.match_score), - band=band, - model_name=row.model, - is_current=True, - ) - self.session.add(entry) - # Flush the INSERT before touching prev/link: with no relationship() - # edge the unit of work emits the UPDATEs first, and both FKs - # (superseded_by_id, inbox.ats_id) reject a pointer to a row that is - # not inserted yet. - await self.session.flush() - if prev is not None: - prev.is_current=False - prev.superseded_by_id=entry.id - self.session.add(prev) - # inbox.ats_id always points at the CURRENT score row. - link.ats_id=entry.id - link.updated_at=datetime.now(timezone.utc) - self.session.add(link) - await self.session.commit() + 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 + 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, + }) async def _sync_upload_ats(self,job,row): - """History row for an upload-sourced score — no inbox application exists, - so inbox_id stays NULL and the supersede chain runs per candidates row - (stable across re-scores: upsert keeps the id for the same job+bytes).""" + """History row for an upload-sourced score — inbox_id stays NULL.""" band=CandidateView._recommendation(row.match_score) or "" - prev=await AtsResults.get_current_for_candidate(self.session,row.id) - entry=AtsResults( - inbox_id=None, - candidate_id=row.id, - job_post_id=job.id, - overall_score=float(row.match_score), - band=band, - model_name=row.model, - is_current=True, - ) - self.session.add(entry) - # Same flush-before-pointing rule as the inbox path: the FK on - # superseded_by_id must see the new row inserted first. - await self.session.flush() - if prev is not None: - prev.is_current=False - prev.superseded_by_id=entry.id - self.session.add(prev) - await self.session.commit() + identity=await AtsResults.resolve_identity(self.session,row.candidate_email,row.id) + 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, + }) class CandidateView: @@ -564,23 +481,14 @@ class CandidateView: # Same bands the frontend uses (Candidates.jsx / seed.js). return "Strong Match" if score>=82 else "Potential Match" if score>=65 else "Weak Match" - async def _scores_by_message(self,message_ids): - """Completed ATS scores (candidates table) per inbox message id, newest first. - - One batched query — the profile list would otherwise pay a query per row. - """ - mids=[m for m in message_ids if m] - if not mids: - return {} - result=await self.session.execute( - select(Candidates) - .where(Candidates.inbox_message_id.in_(mids),Candidates.status=="completed") - .order_by(Candidates.updated_at.desc()) - ) - scores={} - for row in result.scalars().all(): - scores.setdefault(row.inbox_message_id,[]).append(row) - return scores + @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): try: @@ -632,7 +540,19 @@ class CandidateView: job_post=None if manual.job_post_id: job_post=await JobPosts.get_job_post_by_id(self.session,str(manual.job_post_id)) - return serialize_manual_candidate_profile(manual,user,job_post) + 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 return await self.attach_job_posts(rows) except HTTPException: raise @@ -729,7 +649,6 @@ class CandidateView: """ single=not isinstance(data,list) records=[data] if single else list(data or []) - scores=await self._scores_by_message([getattr(r,"message_id",None) for r in records]) payloads=[] wanted=[] @@ -737,10 +656,10 @@ class CandidateView: payload=serialize_candidate_profile(record) payload["job_posts"]=[] payload["assigned_job_post"]=None - scored=scores.get(getattr(record,"message_id",None)) or [] - if scored: - payload["ai_score"]=scored[0].match_score - payload["recommendation"]=self._recommendation(scored[0].match_score) + 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("assigned_job_post_id"): wanted.append(payload["assigned_job_post_id"]) @@ -821,26 +740,46 @@ class CandidateView: ) notes=[serialize_note(r) for r in result.scalars().all()] - # ATS score: join the scoring engine's `candidates` rows onto the profile - # by inbox message. The serializer stubs ai_score/recommendation to None; - # this is where they get real values. Prefer the score against the - # assigned job post, else the most recent completed score. - scores=await self._scores_by_message([getattr(r,"message_id",None) for r in records]) - scored_rows=[row for rows in scores.values() for row in rows] - if scored_rows: - assigned_uid=Candidates._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((r for r in scored_rows if r.job_id==assigned_uid),None) - if chosen is None: - chosen=max(scored_rows,key=lambda r:r.updated_at) - base["ai_score"]=chosen.match_score - base["recommendation"]=self._recommendation(chosen.match_score) - base["matched_keywords"]=list(chosen.matched_keywords or []) - base["missing_keywords"]=list(chosen.missing_keywords or []) - base["summary_critique"]=chosen.summary_critique - base["scored_job_post_id"]=str(chosen.job_id) - base["scored_at"]=chosen.updated_at.isoformat() if chosen.updated_at else None + # 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=[] + if ats_ids: + result=await self.session.execute(select(AtsResults).where(AtsResults.id.in_(ats_ids))) + ats_rows=list(result.scalars().all()) + 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 diff --git a/backend/job/pipeline/serializers.py b/backend/job/pipeline/serializers.py index 030bb84..63bf0f6 100644 --- a/backend/job/pipeline/serializers.py +++ b/backend/job/pipeline/serializers.py @@ -1,7 +1,29 @@ +from inbox.enums import Candidate_application_Status + + +def serialize_pipeline_counts(inbox_counts,manual_counts) -> dict: + by_status={stage.value:0 for stage in Candidate_application_Status} + by_status["UNKNOWN"]=0 + inbox_n=0 + manual_n=0 + for key,n in (inbox_counts or {}).items(): + n=int(n or 0) + inbox_n+=n + bucket=key if key in by_status else "UNKNOWN" + by_status[bucket]+=n + for key,n in (manual_counts or {}).items(): + n=int(n or 0) + manual_n+=n + bucket=key if key in by_status else "UNKNOWN" + by_status[bucket]+=n + return {"by_status":by_status,"inbox":inbox_n,"manual_upload":manual_n} + + def serialize_stage_transition(row) -> dict: return { "id": str(row.id), "inbox_id": row.inbox_id, + "manual_upload_candidate_id": str(row.manual_upload_candidate_id) if row.manual_upload_candidate_id else None, "from_stage": row.from_stage, "to_stage": row.to_stage, "valid_from": row.valid_from.isoformat() if row.valid_from else None, diff --git a/backend/job/pipeline/views.py b/backend/job/pipeline/views.py index 28a9e3e..1d34dec 100644 --- a/backend/job/pipeline/views.py +++ b/backend/job/pipeline/views.py @@ -3,55 +3,92 @@ from sqlalchemy.ext.asyncio import AsyncSession from inbox.enums import Candidate_application_Status from inbox.models import Inbox -from job.candidate.models import ApplicationStageTransitions -from job.pipeline.serializers import serialize_stage_transition - +from job.candidate.models import ApplicationStageTransitions, Manual_UPLOAD_CANDIDATE, _now +from job.pipeline.serializers import serialize_pipeline_counts, serialize_stage_transition +from inbox.plugins import get_ats_score_for_manual_user, get_ats_score_for_user class Pipeline: def __init__(self,session:AsyncSession): self.session=session - async def get_transitions(self,inbox_id=None,transition_id=None): + async def get_all(self,job_post_id=None,limit=None,offset=0): + # limit/offset are per-source, not a merged page: two tables, no common + # order key. limit=200 returns up to 200 inbox AND up to 200 manual rows. + try: + inbox_data=await Inbox.get_all(self.session,job_post_id=job_post_id,limit=limit,offset=offset) + manual_upload_data=await Manual_UPLOAD_CANDIDATE.get_all(self.session,job_post_id=job_post_id,limit=limit,offset=offset) + counts=serialize_pipeline_counts( + await Inbox.count_by_status(self.session,job_post_id=job_post_id), + await Manual_UPLOAD_CANDIDATE.count_by_status(self.session,job_post_id=job_post_id), + ) + return { + "data":{"inbox":inbox_data,"manual_upload":manual_upload_data}, + "counts":counts, + "total":counts["inbox"]+counts["manual_upload"], + } + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + async def get_pipeline_candidates(self,user_id=None,job_post_id=None): + try: + manual_data=await get_ats_score_for_manual_user(self.session,user_id,job_post_id) + inbox_data=await get_ats_score_for_user(self.session,user_id,job_post_id) + data={"manual":manual_data,"inbox":inbox_data} + return data + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + async def get_transitions(self,inbox_id=None,transition_id=None,manual_upload_id=None): if transition_id: row=await ApplicationStageTransitions.get_by_id(self.session,transition_id) if not row: raise HTTPException(status_code=404,detail="Transition not found") return serialize_stage_transition(row) - if inbox_id is None: - raise HTTPException(status_code=400,detail="transition_id or inbox_id is required") - rows=await ApplicationStageTransitions.fetch_by_inbox(self.session,int(inbox_id)) - return [serialize_stage_transition(r) for r in rows] + if inbox_id is not None: + rows=await ApplicationStageTransitions.fetch_by_inbox(self.session,int(inbox_id)) + return [serialize_stage_transition(r) for r in rows] + if manual_upload_id is not None: + rows=await ApplicationStageTransitions.fetch_by_manual(self.session,manual_upload_id) + return [serialize_stage_transition(r) for r in rows] + raise HTTPException(status_code=400,detail="transition_id, inbox_id or manual_upload_id is required") - async def change_stage(self,inbox_id,to_stage,current_user,change_reason=None): + async def change_stage(self,to_stage,current_user,inbox_id=None,manual_upload_id=None,change_reason=None): + if (inbox_id is None)==(manual_upload_id is None): + raise HTTPException(status_code=400,detail="inbox_id or manual_upload_id is required") + try: + stage=Candidate_application_Status(to_stage) + except ValueError: + raise HTTPException(status_code=422,detail="Invalid to_stage") + changed_by=None + if isinstance(current_user,dict) and current_user.get("id"): + changed_by=ApplicationStageTransitions._as_uuid(current_user.get("id")) + if inbox_id is not None: + return await self._change_inbox_stage(inbox_id,stage,changed_by,change_reason) + return await self._change_manual_stage(manual_upload_id,stage,changed_by,change_reason) + + async def _change_inbox_stage(self,inbox_id,stage,changed_by,change_reason): inbox=await Inbox.get_inbox_with_message(self.session,inbox_id) if not inbox: raise HTTPException(status_code=404,detail="Inbox not found") message=inbox.messages if not message: raise HTTPException(status_code=404,detail="Inbox message not found") - try: - stage=Candidate_application_Status(to_stage) - except ValueError: - raise HTTPException(status_code=422,detail="Invalid to_stage") current=message.application_status from_stage=current.value if isinstance(current,Candidate_application_Status) else str(current) if from_stage==stage.value: raise HTTPException(status_code=400,detail="already at stage") - changed_by=None - if isinstance(current_user,dict) and current_user.get("id"): - changed_by=ApplicationStageTransitions._as_uuid(current_user.get("id")) await ApplicationStageTransitions.close_open(self.session,inbox.id,commit=False) - transition_data={ + transition=await ApplicationStageTransitions.insert_transition( + self.session, + { "inbox_id":inbox.id, + "manual_upload_candidate_id":None, "from_stage":from_stage, "to_stage":stage.value, "changed_by":changed_by, "actor_kind":"user", "change_reason":change_reason, - } - transition=await ApplicationStageTransitions.insert_transition( - self.session, - transition_data, + }, commit=False, ) message.application_status=stage @@ -59,6 +96,41 @@ class Pipeline: await self.session.commit() return { "inbox_id":inbox.id, + "manual_upload_id":None, + "application_status":stage.value, + "transition":serialize_stage_transition(transition), + } + + async def _change_manual_stage(self,manual_upload_id,stage,changed_by,change_reason): + row=await Manual_UPLOAD_CANDIDATE.get_by_id(self.session,manual_upload_id) + if not row: + raise HTTPException(status_code=404,detail="Manual upload candidate not found") + from_stage=(row.status or "").strip() or None + if from_stage==stage.value: + raise HTTPException(status_code=400,detail="already at stage") + await ApplicationStageTransitions.close_open( + self.session,manual_upload_candidate_id=row.id,commit=False, + ) + transition=await ApplicationStageTransitions.insert_transition( + self.session, + { + "inbox_id":None, + "manual_upload_candidate_id":row.id, + "from_stage":from_stage, + "to_stage":stage.value, + "changed_by":changed_by, + "actor_kind":"user", + "change_reason":change_reason, + }, + commit=False, + ) + row.status=stage.value + row.updated_at=_now() + self.session.add(row) + await self.session.commit() + return { + "inbox_id":None, + "manual_upload_id":str(row.id), "application_status":stage.value, "transition":serialize_stage_transition(transition), } diff --git a/backend/migrations/manual/002_backfill_inbox_ats.sql b/backend/migrations/manual/002_backfill_inbox_ats.sql index c8d1d4a..4075ad1 100644 --- a/backend/migrations/manual/002_backfill_inbox_ats.sql +++ b/backend/migrations/manual/002_backfill_inbox_ats.sql @@ -8,87 +8,119 @@ -- INSERT skips any application that already has an ats_results row. -- Applied automatically at startup by alembic_setup.run_manual_sql() once the -- schema is at head; recorded in manual_migrations. Safe to re-run by hand. +-- +-- The Candidates.inbox_message_id join is historical: that column is dropped +-- after revision b7d4e8f1a203. Those three statements run only when the column +-- is still present (a DB that has not yet reached that revision). Fresh installs +-- skip them — there is no pre-denorm history to copy. --- Winning score per message, mirroring CandidateView._recommendation precedence: --- the completed score against the ASSIGNED job wins, else the newest completed. -WITH ranked AS ( - SELECT - c.inbox_message_id, - c.job_id, - c.match_score, - c.model, - c.updated_at, - ROW_NUMBER() OVER ( - PARTITION BY c.inbox_message_id - ORDER BY (c.job_id = m.assigned_job_post_id) DESC NULLS LAST, - c.updated_at DESC - ) AS rn - FROM app.candidates c - JOIN app.inbox_messages m ON m.id = c.inbox_message_id - WHERE c.status = 'completed' - AND c.inbox_message_id IS NOT NULL -) -UPDATE app.inbox_messages m -SET ats_score = r.match_score, - ats_band = CASE - WHEN r.match_score >= 82 THEN 'Strong Match' - WHEN r.match_score >= 65 THEN 'Potential Match' - ELSE 'Weak Match' - END -FROM ranked r -WHERE r.rn = 1 - AND m.id = r.inbox_message_id - AND m.ats_score IS NULL; +DO $backfill$ +BEGIN + IF EXISTS ( + SELECT 1 FROM information_schema.columns + WHERE table_schema = 'app' + AND table_name = 'candidates' + AND column_name = 'inbox_message_id' + ) THEN + -- Winning score per message, mirroring CandidateView._recommendation precedence: + -- the completed score against the ASSIGNED job wins, else the newest completed. + EXECUTE $sql$ + WITH ranked AS ( + SELECT + c.inbox_message_id, + c.job_id, + c.match_score, + c.model, + c.updated_at, + ROW_NUMBER() OVER ( + PARTITION BY c.inbox_message_id + ORDER BY (c.job_id = m.assigned_job_post_id) DESC NULLS LAST, + c.updated_at DESC + ) AS rn + FROM app.candidates c + JOIN app.inbox_messages m ON m.id = c.inbox_message_id + WHERE c.status = 'completed' + AND c.inbox_message_id IS NOT NULL + ) + UPDATE app.inbox_messages m + SET ats_score = r.match_score, + ats_band = CASE + WHEN r.match_score >= 82 THEN 'Strong Match' + WHEN r.match_score >= 65 THEN 'Potential Match' + ELSE 'Weak Match' + END + FROM ranked r + WHERE r.rn = 1 + AND m.id = r.inbox_message_id + AND m.ats_score IS NULL + $sql$; --- One current ats_results row per already-scored application. computed_at takes --- the candidates row's timestamp so the history reflects when the score happened. --- inbox can hold several join rows per message; DISTINCT ON keeps the newest. -WITH ranked AS ( - SELECT - c.inbox_message_id, - c.job_id, - c.match_score, - c.model, - c.updated_at, - ROW_NUMBER() OVER ( - PARTITION BY c.inbox_message_id - ORDER BY (c.job_id = m.assigned_job_post_id) DESC NULLS LAST, - c.updated_at DESC - ) AS rn - FROM app.candidates c - JOIN app.inbox_messages m ON m.id = c.inbox_message_id - WHERE c.status = 'completed' - AND c.inbox_message_id IS NOT NULL -), -links AS ( - SELECT DISTINCT ON (message_id) message_id, id AS inbox_id - FROM app.inbox - ORDER BY message_id, created_at DESC -) -INSERT INTO app.ats_results - (id, inbox_id, job_post_id, overall_score, band, is_current, - superseded_by_id, model_name, computed_at, created_at) -SELECT - gen_random_uuid(), - l.inbox_id, - r.job_id, - r.match_score, - CASE - WHEN r.match_score >= 82 THEN 'Strong Match' - WHEN r.match_score >= 65 THEN 'Potential Match' - ELSE 'Weak Match' - END, - true, - NULL, - r.model, - r.updated_at, - NOW() -FROM ranked r -JOIN links l ON l.message_id = r.inbox_message_id -WHERE r.rn = 1 - AND NOT EXISTS ( - SELECT 1 FROM app.ats_results a WHERE a.inbox_id = l.inbox_id - ); + -- One current ats_results row per already-scored application. computed_at takes + -- the candidates row's timestamp so the history reflects when the score happened. + -- inbox can hold several join rows per message; DISTINCT ON keeps the newest. + EXECUTE $sql$ + WITH ranked AS ( + SELECT + c.inbox_message_id, + c.job_id, + c.match_score, + c.model, + c.updated_at, + ROW_NUMBER() OVER ( + PARTITION BY c.inbox_message_id + ORDER BY (c.job_id = m.assigned_job_post_id) DESC NULLS LAST, + c.updated_at DESC + ) AS rn + FROM app.candidates c + JOIN app.inbox_messages m ON m.id = c.inbox_message_id + WHERE c.status = 'completed' + AND c.inbox_message_id IS NOT NULL + ), + links AS ( + SELECT DISTINCT ON (message_id) message_id, id AS inbox_id + FROM app.inbox + ORDER BY message_id, created_at DESC + ) + INSERT INTO app.ats_results + (id, inbox_id, job_post_id, overall_score, band, is_current, + superseded_by_id, model_name, computed_at, created_at) + SELECT + gen_random_uuid(), + l.inbox_id, + r.job_id, + r.match_score, + CASE + WHEN r.match_score >= 82 THEN 'Strong Match' + WHEN r.match_score >= 65 THEN 'Potential Match' + ELSE 'Weak Match' + END, + true, + NULL, + r.model, + r.updated_at, + NOW() + FROM ranked r + JOIN links l ON l.message_id = r.inbox_message_id + WHERE r.rn = 1 + AND NOT EXISTS ( + SELECT 1 FROM app.ats_results a WHERE a.inbox_id = l.inbox_id + ) + $sql$; + + -- Older inbox history rows predate candidate_id; link the current ones back to + -- the candidates row that produced them. + EXECUTE $sql$ + UPDATE app.ats_results a + SET candidate_id = c.id + FROM app.inbox i + JOIN app.candidates c ON c.inbox_message_id = i.message_id AND c.status = 'completed' + WHERE a.inbox_id = i.id + AND a.candidate_id IS NULL + AND c.job_id = a.job_post_id + $sql$; + END IF; +END +$backfill$; -- inbox.ats_id -> the CURRENT ats_results row for that application, so the -- score is one direct id join away. Newest current row wins if several exist. @@ -131,13 +163,3 @@ WHERE c.status = 'completed' AND NOT EXISTS ( SELECT 1 FROM app.ats_results a WHERE a.candidate_id = c.id ); - --- Older inbox history rows predate candidate_id; link the current ones back to --- the candidates row that produced them. -UPDATE app.ats_results a -SET candidate_id = c.id -FROM app.inbox i -JOIN app.candidates c ON c.inbox_message_id = i.message_id AND c.status = 'completed' -WHERE a.inbox_id = i.id - AND a.candidate_id IS NULL - AND c.job_id = a.job_post_id; diff --git a/backend/users/models.py b/backend/users/models.py index 0ecae32..e7f8c54 100644 --- a/backend/users/models.py +++ b/backend/users/models.py @@ -57,7 +57,6 @@ class Users(SQLModel, table=True): ) password: str - created_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) updated_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) is_active: bool = Field(default=False) diff --git a/frontend/dist/index.html b/frontend/dist/index.html index db30cad..2c57fc8 100644 --- a/frontend/dist/index.html +++ b/frontend/dist/index.html @@ -23,7 +23,7 @@ - +
diff --git a/frontend/src/api/candidates.js b/frontend/src/api/candidates.js index 1b37bf0..eae3e04 100644 --- a/frontend/src/api/candidates.js +++ b/frontend/src/api/candidates.js @@ -83,7 +83,6 @@ export function toCandidateView(row) { errorCode: row.error_code ?? null, errorMessage: row.error_message ?? null, applied: row.created_at ? new Date(row.created_at) : null, - inboxMessageId: row.inbox_message_id ?? null, } } @@ -115,7 +114,7 @@ export function listCandidateUsers({ roleId = 8, top = 500, skip = 0 } = {}) { * * A user account carries identity only. Everything the ATS produces * (score, matched skills, critique, the job it was scored against) lives in the - * `candidates` table keyed by inbox_message_id, with no user_id to join on, so + * `candidates` table keyed by job_id + content hash, with no user_id to join on, so * those fields are null here by construction rather than by omission. */ export function toCandidateUserView(row) { @@ -141,7 +140,6 @@ export function toCandidateUserView(row) { scoringStatus: null, errorCode: null, errorMessage: null, - inboxMessageId: null, } } diff --git a/frontend/src/api/pipeline.js b/frontend/src/api/pipeline.js index c1c0dab..4f377b0 100644 --- a/frontend/src/api/pipeline.js +++ b/frontend/src/api/pipeline.js @@ -1,17 +1,13 @@ /* ============================================================ pipeline.js — the kanban board's endpoints (backend/job/app.py). - The board is stitched from two modules, because there is no pipeline-specific - READ endpoint: + The board is one read: GET /pipeline/candidates/fetch returns inbox + + manual_upload cards plus per-status counts. Dropping a card fires + PATCH /candidate/stage with `inbox_id` or `manual_upload_id`. - - rows come from GET /candidate/fetch (api/candidates.js `list`), the - inbox -> users -> roles join, which is the only list payload carrying BOTH - `application_status` (the stage) and `inbox_id` (what the write below needs); - - the write is PATCH /candidate/stage, here. - - Stage lives on inbox_messages.application_status and the transition history in - application_stage_transitions; the server closes the open interval and opens a - new one in the same commit, so the board never has to touch history itself. + Inbox stage lives on inbox_messages.application_status; manual stage lives + on manual_upload_candidate.status (same Candidate_application_Status values). + History is application_stage_transitions in both cases. ============================================================ */ import { request } from '../lib/apiClient' @@ -60,63 +56,107 @@ export const STATUS_FROM_STAGE = { /** * Move one application to another stage. Requires pipeline.edit. * - * `inboxId` is the INTEGER inbox.id — the row the candidate profile returns as - * `inbox_id`, not the inbox_messages uuid the Inbox screen calls `id`; the route - * runs int() on it and 404s on anything else. - * - * The server rejects a no-op move with 400 ("already at stage"), so callers must - * not fire on a drop into the card's current column. + * Inbox cards send `inboxId` (INTEGER inbox.id). Manual-upload cards send + * `manualUploadId` (manual_upload_candidate.id UUID). The server 400s if both + * or neither are present, and 400s a no-op move ("already at stage"). */ -export function changeStage({ inboxId, toStage, changeReason }) { +export function changeStage({ inboxId, manualUploadId, toStage, changeReason }) { return request('/candidate/stage', { method: 'PATCH', - body: { inbox_id: inboxId, to_stage: toStage, change_reason: changeReason ?? null }, + body: { + ...(inboxId != null ? { inbox_id: inboxId } : {}), + ...(manualUploadId != null ? { manual_upload_id: manualUploadId } : {}), + to_stage: toStage, + change_reason: changeReason ?? null, + }, }) } +/** + * Inbox + manual-upload applications for the board — GET /pipeline/candidates/fetch + * (pipeline.view). Envelope is `{ data: { inbox, manual_upload }, counts, total }`. + * `jobId === ''` (All Jobs) is dropped by buildUrl and sends no filter. + */ +export function listApplications({ jobId, limit, offset } = {}) { + return request('/pipeline/candidates/fetch', { + params: { job_post_id: jobId, limit, offset }, + }) +} + +/** + * Fold the 11 status counts into the 7 board columns. Unmapped keys (UNKNOWN) + * land in Applied, same as STAGE_FROM_STATUS's card fallback. + */ +export function toStageCounts(byStatus) { + const counts = Object.fromEntries(Object.keys(STATUS_FROM_STAGE).map((name) => [name, 0])) + for (const [status, n] of Object.entries(byStatus || {})) { + const stage = STAGE_FROM_STATUS[status] ?? 'Applied' + counts[stage] = (counts[stage] ?? 0) + (n || 0) + } + return counts +} + /** * Stage history for one application — GET /pipeline/transitions/fetch * (pipeline.view). Rows are valid-time intervals: `valid_to` null is the stage * the candidate is in now. The board itself does not render history; this is the * feed behind a stage timeline on the profile. * - * One of inboxId / transitionId is required — the route 400s with neither. + * One of inboxId / manualUploadId / transitionId is required — the route 400s + * with none of them. */ -export function listTransitions({ inboxId, transitionId } = {}) { +export function listTransitions({ inboxId, manualUploadId, transitionId } = {}) { return request('/pipeline/transitions/fetch', { - params: { inbox_id: inboxId, transition_id: transitionId }, + params: { inbox_id: inboxId, manual_upload_id: manualUploadId, transition_id: transitionId }, }) } -/** - * Candidate-profile row -> one kanban card. - * - * `id` is the inbox id, not the user id: the board is one card per APPLICATION - * and `inbox` holds one row per (user, message), so a candidate who mailed us - * three times legitimately occupies three cards with three independent stages. - * `userId` rides along for the deep link into the profile. - * - * Skills are absent by construction — the list payload carries ai_score but not - * matched_keywords (job/candidate/views.py::attach_job_posts sets only the - * score), so the card drops its tag row rather than rendering three blanks. - */ -export function toBoardCard(row) { - const jobTitle = row.job_title ?? row.assigned_job_post?.title ?? null +function sourceFields(row, kind) { + if (kind === 'manual') { + return { + id: `manual:${row.id}`, + inboxId: null, + manualUploadId: row.id, + jobId: row.job_post_id ?? null, + jobTitle: row.title ?? null, + currentTitle: row.current_position || null, + currentCompany: row.current_company || null, + } + } return { id: row.inbox_id, inboxId: row.inbox_id, + manualUploadId: null, + jobId: row.assigned_job_post_id ?? null, + jobTitle: row.title ?? null, + currentTitle: row.current_title || null, + currentCompany: row.current_employment || null, + } +} + +/** + * Pipeline inbox or manual-upload row -> one kanban card. + * + * Inbox `id` is the inbox id, not the user id: the board is one card per + * APPLICATION. `userId` rides along for the deep link into the profile. + */ +export function toBoardCard(row, kind = 'inbox') { + const src = sourceFields(row, kind) + return { + ...src, userId: row.user_id ?? null, name: row.name || row.email || 'Unknown', email: row.email ?? null, stage: STAGE_FROM_STATUS[row.application_status] ?? 'Applied', status: row.application_status ?? null, - jobId: row.assigned_job_post_id ?? null, - jobTitle, - currentTitle: row.current_title || null, - currentCompany: row.current_employment || null, experience: row.experience || null, - aiScore: row.ai_score ?? null, - recommendation: row.recommendation ?? null, + aiScore: row.ats_result?.overall_score ?? null, + recommendation: row.ats_result?.band ?? null, applied: row.created_at ? new Date(row.created_at) : null, } } + +/** GET /pipeline/candidates/fetch `manual_upload` row -> one kanban card. */ +export function toManualBoardCard(row) { + return toBoardCard(row, 'manual') +} diff --git a/frontend/src/screens/Pipeline.jsx b/frontend/src/screens/Pipeline.jsx index 9c75bb5..0f7ae5e 100644 --- a/frontend/src/screens/Pipeline.jsx +++ b/frontend/src/screens/Pipeline.jsx @@ -1,31 +1,28 @@ /* ============================================================ Pipeline — the kanban board, on live backend data. - Cards come from GET /candidate/fetch (the inbox -> users -> roles join), the - only list payload that carries the stage (`application_status`) together with - the `inbox_id` that PATCH /candidate/stage writes against. Dropping a card - fires that PATCH; the server closes the open application_stage_transitions - interval and opens a new one in the same commit. + Cards and column counts come from GET /pipeline/candidates/fetch. Dropping a + card fires PATCH /candidate/stage with `inbox_id` or `manual_upload_id`. The + server closes the open application_stage_transitions interval and opens a + new one in the same commit. - The job filter reads live posts from GET /job/fetch and matches on - `assigned_job_post_id`, so an application nobody has assigned to a post shows - under All Jobs only. + The job filter is server-side (`job_post_id`). Applications with no assigned + post are omitted by the list query, so they do not appear under All Jobs. - The card's skill tags are gone: the list payload has ai_score but no + The card's skill tags are gone: the list payload has ats_result but no matched_keywords, and the Candidates screen set the precedent that a column with no source is dropped rather than rendered as blanks. ============================================================ */ import { useMemo, useState } from 'react' import { useNavigate } from 'react-router-dom' -import { useMutation, useQuery, useQueryClient } from '@tanstack/react-query' +import { keepPreviousData, useMutation, useQuery, useQueryClient } from '@tanstack/react-query' import { Avatar, EmptyState, Icon, ScoreChip } from '../ui/primitives' import { useToast } from '../ui/Toast' import { useAuth } from '../auth/AuthContext' import { qk } from '../lib/queryKeys' import { friendlyAuthError } from '../lib/errors' -import * as candidatesApi from '../api/candidates' import * as jobPostsApi from '../api/jobPosts' import * as pipelineApi from '../api/pipeline' @@ -40,18 +37,21 @@ export const KANBAN_STAGES = [ { name: 'Rejected', color: 'var(--stage-7)' }, ] -/* /candidate/fetch pages with limit/offset and has no job filter, so the board - pulls one page and filters client-side. Rows past this are not on the board — - the header says so rather than silently showing a partial pipeline. */ const BOARD_LIMIT = 200 -const BOARD_KEY = qk.pipeline.board({ limit: BOARD_LIMIT }) - const JOB_LIMIT = 100 -async function fetchBoard() { - const res = await candidatesApi.list({ limit: BOARD_LIMIT }) - const rows = candidatesApi.toRows(res) - return { cards: rows.map(pipelineApi.toBoardCard), total: res?.total ?? rows.length } +async function fetchBoard(jobId) { + const res = await pipelineApi.listApplications({ jobId, limit: BOARD_LIMIT }) + const inbox = Array.isArray(res?.data?.inbox) ? res.data.inbox : [] + const manuals = Array.isArray(res?.data?.manual_upload) ? res.data.manual_upload : [] + return { + cards: [ + ...inbox.map((row) => pipelineApi.toBoardCard(row)), + ...manuals.map((row) => pipelineApi.toManualBoardCard(row)), + ], + total: res?.total ?? 0, + stageCounts: pipelineApi.toStageCounts(res?.counts?.by_status), + } } async function fetchJobs() { @@ -74,53 +74,67 @@ export default function Pipeline() { const navigate = useNavigate() const qc = useQueryClient() - const board = useQuery({ queryKey: BOARD_KEY, queryFn: fetchBoard }) - const { data: jobs = [] } = useQuery({ queryKey: qk.jobPosts.list({ top: JOB_LIMIT }), queryFn: fetchJobs }) - const [jobId, setJobId] = useState('') const [draggingId, setDraggingId] = useState(null) const [overStage, setOverStage] = useState(null) + const boardKey = useMemo( + () => qk.pipeline.board({ limit: BOARD_LIMIT, jobId: jobId || null }), + [jobId], + ) + + const board = useQuery({ + queryKey: boardKey, + queryFn: () => fetchBoard(jobId), + placeholderData: keepPreviousData, + }) + const { data: jobs = [] } = useQuery({ queryKey: qk.jobPosts.list({ top: JOB_LIMIT }), queryFn: fetchJobs }) + /* The route is behind pipeline.view, but the WRITE needs pipeline.edit — a viewer gets a read-only board instead of drags that 403 on drop. */ const canEdit = can('pipeline.edit') const candidates = board.data?.cards ?? [] const total = board.data?.total ?? 0 - - const list = useMemo( - () => (jobId ? candidates.filter((c) => c.jobId === jobId) : candidates), - [candidates, jobId], - ) + const stageCounts = board.data?.stageCounts ?? {} const byStage = useMemo(() => { const map = Object.fromEntries(KANBAN_STAGES.map((s) => [s.name, []])) - for (const c of list) if (map[c.stage]) map[c.stage].push(c) + for (const c of candidates) if (map[c.stage]) map[c.stage].push(c) return map - }, [list]) + }, [candidates]) /* Optimistic: a drag that only repaints after the round trip reads as a failed drop. The card snaps back on error and the server's own value wins on the - refetch in onSettled. */ + refetch in onSettled. Counts move with the card so badges do not lag. */ const move = useMutation({ mutationFn: ({ card, stage }) => pipelineApi.changeStage({ inboxId: card.inboxId, + manualUploadId: card.manualUploadId, toStage: pipelineApi.STATUS_FROM_STAGE[stage], }), onMutate: async ({ card, stage }) => { - await qc.cancelQueries({ queryKey: BOARD_KEY }) - const previous = qc.getQueryData(BOARD_KEY) - qc.setQueryData(BOARD_KEY, (old) => - old && { + await qc.cancelQueries({ queryKey: boardKey }) + const previous = qc.getQueryData(boardKey) + qc.setQueryData(boardKey, (old) => { + if (!old) return old + const from = card.stage + const nextCounts = { ...old.stageCounts } + if (from && from !== stage) { + nextCounts[from] = Math.max(0, (nextCounts[from] ?? 0) - 1) + nextCounts[stage] = (nextCounts[stage] ?? 0) + 1 + } + return { ...old, - cards: old.cards.map((c) => (c.inboxId === card.inboxId ? { ...c, stage } : c)), - }, - ) + cards: old.cards.map((c) => (c.id === card.id ? { ...c, stage } : c)), + stageCounts: nextCounts, + } + }) return { previous } }, onError: (err, _vars, ctx) => { - if (ctx?.previous) qc.setQueryData(BOARD_KEY, ctx.previous) + if (ctx?.previous) qc.setQueryData(boardKey, ctx.previous) toast(friendlyAuthError(err, 'Could not move the candidate.'), 'error') }, onSuccess: (_data, { card, stage }) => toast(`${card.name} moved to ${stage}`, 'success'), @@ -140,7 +154,7 @@ export default function Pipeline() { const cand = candidates.find((c) => c.id === id) // A no-op move is a 400 server-side ("already at stage"), so it never leaves. if (!cand || cand.stage === stage) return - if (cand.inboxId == null) { + if (cand.inboxId == null && cand.manualUploadId == null) { toast(`${cand.name} has no application to move`, 'warning') return } @@ -190,7 +204,7 @@ export default function Pipeline() {