137 lines
6.4 KiB
Python
137 lines
6.4 KiB
Python
from fastapi import HTTPException
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
from inbox.enums import Candidate_application_Status
|
|
from inbox.models import Inbox
|
|
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_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 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,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")
|
|
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")
|
|
await ApplicationStageTransitions.close_open(self.session,inbox.id,commit=False)
|
|
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,
|
|
},
|
|
commit=False,
|
|
)
|
|
message.application_status=stage
|
|
self.session.add(message)
|
|
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),
|
|
}
|