HR-ATS-Portal/backend/offer/views.py

560 lines
23 KiB
Python

import uuid
from datetime import datetime,timezone
import httpx
from fastapi import HTTPException
from sqlalchemy.ext.asyncio import AsyncSession
from sqlmodel import select
from g_sheet.models import FormData
from inbox.models import Inbox
from job.candidate.models import Manual_UPLOAD_CANDIDATE
from job.candidate.views import owned_job_ids_for_candidate_scope
from job.history.enums import HistoryEvent
from job.history.views import HistoryRecorder
from job.job_post.models import JobPosts
from job.pipeline.views import Pipeline
from notifications.models import Notifications
from offer.models import Offers,OfferStatusHistory
from offer.plugins import (
email_key,
is_interview_plus,
non_validation_values,
parse_offer_datetime,
render_offer_email,
send_offer_mail,
stage_value,
)
from offer.serializers import serialize_offer,serialize_offer_candidate
from users.models import Users
from users.permissions import sees_all_offers
SOURCE_RANK={"inbox":0,"manual":1,"form":2}
MAIL_FAIL_DETAIL="Failed: the offer could not be sent"
def _as_uuid(value):
if value in (None,""):
return None
try:
return uuid.UUID(str(value))
except (TypeError,ValueError):
return None
def _user_id(current_user):
if not current_user or not current_user.get("id"):
raise HTTPException(status_code=401,detail="Not authenticated")
uid=_as_uuid(current_user["id"])
if uid is None:
raise HTTPException(status_code=401,detail="Invalid user id")
return uid
def _comp_fields(payload):
fields={}
for key in non_validation_values():
if key not in payload:
continue
value=payload[key]
if key in ("start_date","expiry_date"):
value=parse_offer_datetime(value)
fields[key]=value
return fields
class Offer:
def __init__(self,session:AsyncSession):
self.session=session
async def _hydrate_offers(self,rows):
ids=[]
for row in rows:
ids.append(row.candidate_user_id)
ids.append(row.created_by)
names=await Users.names_by_ids(self.session,ids)
return [
serialize_offer(
row,
candidate_name=names.get(str(row.candidate_user_id)),
created_by_name=names.get(str(row.created_by)),
)
for row in rows
]
async def get_offers(self,offer_id=None,status=None,inbox_id=None,job_post_id=None,top=None,skip=0):
if offer_id is not None:
row=await Offers.get_offer_by_id(self.session,offer_id)
if not row:
raise HTTPException(status_code=404,detail="Offer not found")
return (await self._hydrate_offers([row]))[0],1
rows,total=await Offers.fetch_offers(
self.session,
status=status,
inbox_id=inbox_id,
job_post_id=job_post_id,
top=top,
skip=skip or 0,
)
return await self._hydrate_offers(rows),total
async def _offer_job_ids(self,current_user):
if sees_all_offers(current_user):
return None
return await owned_job_ids_for_candidate_scope(self.session,current_user)
async def _assert_offer_job(self,current_user,job_post_id):
if sees_all_offers(current_user):
return
owned=await owned_job_ids_for_candidate_scope(self.session,current_user)
owned=set(owned or [])
jid=_as_uuid(job_post_id)
if jid is None or jid not in owned:
raise HTTPException(status_code=403,detail="This offer is outside your assigned jobs")
async def list_candidates(self,current_user,search=None,top=None,skip=0):
owned=await self._offer_job_ids(current_user)
if owned is not None and not owned:
return [],0
inbox_rows=await Inbox.get_all(self.session,job_post_ids=owned)
manual_rows=await Manual_UPLOAD_CANDIDATE.get_all(self.session,job_post_ids=owned)
form_rows=await FormData.list_for_offer_picker(
self.session,job_post_ids=owned,search=None,
)
form_by_manual=await FormData.form_ids_by_manual_ids(
self.session,[r.get("id") for r in manual_rows],
)
merged={}
for row in inbox_rows:
if row.get("is_duplicate"):
continue
if not is_interview_plus(row.get("application_status")):
continue
item={
"inbox_id":row.get("inbox_id"),
"manual_upload_candidate_id":None,
"form_data_id":None,
"source":"inbox",
"user_id":row.get("user_id"),
"name":row.get("name"),
"email":email_key(row.get("email")),
"job_post_id":row.get("assigned_job_post_id"),
"job_title":row.get("title"),
"application_status":stage_value(row.get("application_status")),
}
self._merge_candidate(merged,item)
for row in manual_rows:
if not is_interview_plus(row.get("application_status")):
continue
mid=row.get("id")
item={
"inbox_id":None,
"manual_upload_candidate_id":mid,
"form_data_id":form_by_manual.get(str(mid)) if mid else None,
"source":"manual",
"user_id":row.get("user_id"),
"name":row.get("name") or None,
"email":email_key(row.get("email") or row.get("candidate_email")),
"job_post_id":row.get("job_post_id"),
"job_title":row.get("title"),
"application_status":stage_value(row.get("application_status")),
}
self._merge_candidate(merged,item)
for row in form_rows:
item={
"inbox_id":None,
"manual_upload_candidate_id":None,
"form_data_id":row.get("form_data_id"),
"source":"form",
"user_id":row.get("user_id"),
"name":row.get("name"),
"email":email_key(row.get("email")),
"job_post_id":row.get("job_post_id"),
"job_title":row.get("job_title"),
"application_status":stage_value(row.get("application_status")) or "PENDING",
}
self._merge_candidate(merged,item)
items=list(merged.values())
needle=(search or "").strip().lower()
if needle:
items=[
r for r in items
if needle in (r.get("name") or "").lower() or needle in (r.get("email") or "")
]
items.sort(key=lambda r:(r.get("name") or r.get("email") or "").lower())
total=len(items)
start=int(skip or 0)
if start:
items=items[start:]
if top is not None:
items=items[:int(top)]
return [serialize_offer_candidate(r) for r in items],total
def _merge_candidate(self,merged,item):
job_id=item.get("job_post_id")
if not job_id:
return
email=item.get("email") or ""
if email:
key=(email,str(job_id))
else:
source=item.get("source") or "row"
raw=item.get("inbox_id") or item.get("manual_upload_candidate_id") or item.get("form_data_id")
key=(f"noid:{source}:{raw}",str(job_id))
existing=merged.get(key)
if existing is None:
merged[key]=item
return
if SOURCE_RANK.get(item.get("source"),9)<SOURCE_RANK.get(existing.get("source"),9):
if not item.get("form_data_id"):
item["form_data_id"]=existing.get("form_data_id")
if not item.get("name"):
item["name"]=existing.get("name")
merged[key]=item
return
if not existing.get("form_data_id"):
existing["form_data_id"]=item.get("form_data_id")
if not existing.get("name"):
existing["name"]=item.get("name")
async def create_offer(self,payload,current_user):
if not payload.get("inbox_id"):
raise HTTPException(status_code=422,detail="inbox_id is required")
job_post_id=_as_uuid(payload.get("job_post_id"))
if job_post_id is None:
raise HTTPException(status_code=422,detail="job_post_id is required")
candidate_user_id=_as_uuid(payload.get("candidate_user_id"))
if candidate_user_id is None:
raise HTTPException(status_code=422,detail="candidate_user_id is required")
created_by=_user_id(current_user)
status=payload.get("status") or "draft"
fields={
"inbox_id": int(payload["inbox_id"]),
"job_post_id": job_post_id,
"candidate_user_id": candidate_user_id,
"created_by": created_by,
"status": status,
}
for key in non_validation_values():
if key in payload and payload[key] is not None:
fields[key]=payload[key]
row=await Offers.insert_offer(self.session,fields)
history_data={
"offer_id": row.id,
"from_status": None,
"to_status": status,
"changed_by": created_by,
"actor_kind": "user",
"change_reason": payload.get("change_reason"),
}
await OfferStatusHistory.insert_history(self.session,history_data)
return serialize_offer(row)
async def update_offer(self,offer_id,payload,current_user):
row=await Offers.get_offer_by_id(self.session,offer_id)
if not row:
raise HTTPException(status_code=404,detail="Offer not found")
changed_by=_user_id(current_user)
fields={}
for key in (
"status","base_salary","currency","salary_period","signing_bonus","annual_bonus_pct",
"equity_units","equity_instrument","start_date","expiry_date","sent_at",
"responded_at","closed_at","issued_by","inbox_id","job_post_id","candidate_user_id",
"cadre","gross_salary_in_words","subsidized_services","probation_period",
"notice_period","work_location","work_timings",
):
if key not in payload:
continue
value=payload[key]
if key in ("job_post_id","candidate_user_id","issued_by") and value is not None:
value=_as_uuid(value)
if value is None:
raise HTTPException(status_code=422,detail=f"Invalid {key}")
fields[key]=value
if not fields:
raise HTTPException(status_code=400,detail="No fields to update")
new_status=fields.get("status")
if new_status is not None and new_status!=row.status:
await OfferStatusHistory.insert_history(self.session,{
"offer_id": row.id,
"from_status": row.status,
"to_status": new_status,
"changed_by": changed_by,
"actor_kind": "user",
"change_reason": payload.get("change_reason"),
},commit=False)
updated=await Offers.update_offer(self.session,offer_id,fields)
if not updated:
raise HTTPException(status_code=404,detail="Offer not found")
return serialize_offer(updated)
async def issue_offer(self,offer_id,current_user):
row=await Offers.get_offer_by_id(self.session,offer_id)
if not row:
raise HTTPException(status_code=404,detail="Offer not found")
issued_by=_user_id(current_user)
now=datetime.now(timezone.utc)
from_status=row.status
fields={"issued_by": issued_by,"sent_at": now}
to_status=from_status
if from_status=="draft":
fields["status"]="sent"
to_status="sent"
updated=await Offers.update_offer(self.session,offer_id,fields)
if not updated:
raise HTTPException(status_code=404,detail="Offer not found")
history_data={
"offer_id": updated.id,
"from_status": from_status,
"to_status": to_status,
"changed_by": issued_by,
"actor_kind": "user",
"change_reason": "issued",
}
await OfferStatusHistory.insert_history(self.session,history_data)
return serialize_offer(updated)
async def send_offer(self,payload,current_user):
created_by=_user_id(current_user)
app=await self._resolve_application(payload)
await self._assert_offer_job(current_user,app["job_post_id"])
candidate_user_id=app["candidate_user_id"]
job_post_id=app["job_post_id"]
if payload.get("base_salary") in (None,""):
raise HTTPException(status_code=422,detail="base_salary is required")
open_row=await Offers.get_open_for_candidate_job(self.session,candidate_user_id,job_post_id)
retry_id=_as_uuid(payload.get("offer_id"))
if open_row and (retry_id is None or open_row.id!=retry_id):
raise HTTPException(status_code=409,detail="An offer is already in progress for this candidate and job")
fields=_comp_fields(payload)
fields.update({
"inbox_id": app.get("inbox_id"),
"manual_upload_candidate_id": app.get("manual_upload_id"),
"form_data_id": app.get("form_data_id"),
"job_post_id": job_post_id,
"candidate_user_id": candidate_user_id,
"status": "failed",
})
if fields.get("equity_units") in (None,0):
fields["equity_instrument"]=None
row=None
if retry_id is not None:
row=await Offers.get_offer_by_id(self.session,retry_id)
if not row:
raise HTTPException(status_code=404,detail="Offer not found")
if row.status in ("sent","negotiating"):
raise HTTPException(status_code=409,detail="This offer has already been sent")
row=await Offers.update_offer(self.session,retry_id,fields)
else:
failed=await Offers.get_failed_for_candidate_job(self.session,candidate_user_id,job_post_id)
if failed:
row=await Offers.update_offer(self.session,failed.id,fields)
else:
fields["created_by"]=created_by
row=await Offers.insert_offer(self.session,fields)
await OfferStatusHistory.insert_history(self.session,{
"offer_id": row.id,
"from_status": None,
"to_status": "failed",
"changed_by": created_by,
"actor_kind": "user",
"change_reason": payload.get("change_reason"),
})
name,email,job_title=await self._offer_mail_context(row,app)
subject,html=render_offer_email(
name,row.base_salary,row.currency,row.salary_period,
annual_bonus_pct=row.annual_bonus_pct,
signing_bonus=row.signing_bonus,
equity_units=row.equity_units,
equity_instrument=row.equity_instrument,
)
try:
if not email:
raise RuntimeError("candidate has no email")
await send_offer_mail(email,subject,html)
except (httpx.HTTPError,RuntimeError) as e:
await self._notify_send_failed(created_by,name,job_title,row)
raise HTTPException(status_code=502,detail=MAIL_FAIL_DETAIL) from e
now=datetime.now(timezone.utc)
from_status=row.status
updated=await Offers.update_offer(self.session,row.id,{
"status": "sent",
"sent_at": now,
"issued_by": created_by,
})
await OfferStatusHistory.insert_history(self.session,{
"offer_id": updated.id,
"from_status": from_status,
"to_status": "sent",
"changed_by": created_by,
"actor_kind": "user",
"change_reason": "sent",
})
await self._advance_pipeline(app,current_user)
await HistoryRecorder(self.session).record(
HistoryEvent.OFFER_SENT.value,
current_user=current_user,
user_id=candidate_user_id,
inbox_id=app.get("inbox_id"),
manual_upload_candidate_id=app.get("manual_upload_id"),
entity_type="offer",
entity_id=updated.id,
to_value="OFFER",
description=f"Offer sent to {name}",
commit=True,
)
return (await self._hydrate_offers([updated]))[0]
async def _resolve_application(self,payload):
inbox_id=payload.get("inbox_id")
manual_id=_as_uuid(payload.get("manual_upload_candidate_id"))
form_id=_as_uuid(payload.get("form_data_id"))
requested_job=_as_uuid(payload.get("job_post_id"))
present=sum(1 for v in (inbox_id not in (None,""), manual_id is not None, form_id is not None) if v)
if present!=1:
raise HTTPException(
status_code=422,
detail="Exactly one of inbox_id, manual_upload_candidate_id or form_data_id is required",
)
if inbox_id not in (None,""):
return await self._resolve_inbox(inbox_id,requested_job)
if form_id is not None:
return await self._resolve_form(form_id,requested_job)
return await self._resolve_manual(manual_id,requested_job)
async def _resolve_inbox(self,inbox_id,requested_job):
inbox=await Inbox.get_inbox_with_message(self.session,inbox_id)
if not inbox or not inbox.messages:
raise HTTPException(status_code=404,detail="Inbox not found")
message=inbox.messages
job_id=message.assigned_job_post_id
if job_id is None:
raise HTTPException(status_code=422,detail="Candidate is not assigned to a job")
if requested_job is not None and job_id!=requested_job:
raise HTTPException(status_code=422,detail="job_post_id does not match the application")
if not is_interview_plus(message.application_status):
raise HTTPException(status_code=422,detail="Candidate is not at interview stage or later")
if not inbox.user_id:
raise HTTPException(status_code=422,detail="candidate_user_id is required")
return {
"inbox_id": inbox.id,
"manual_upload_id": None,
"form_data_id": None,
"job_post_id": job_id,
"candidate_user_id": inbox.user_id,
}
async def _resolve_manual(self,manual_id,requested_job,require_interview=True):
row=await Manual_UPLOAD_CANDIDATE.get_by_id(self.session,manual_id)
if not row:
raise HTTPException(status_code=404,detail="Manual upload candidate not found")
if not row.job_post_id:
raise HTTPException(status_code=422,detail="Candidate is not assigned to a job")
if requested_job is not None and row.job_post_id!=requested_job:
raise HTTPException(status_code=422,detail="job_post_id does not match the application")
if require_interview and not is_interview_plus(row.status):
raise HTTPException(status_code=422,detail="Candidate is not at interview stage or later")
if not row.user_id:
raise HTTPException(status_code=422,detail="candidate_user_id is required")
form_ids=await FormData.form_ids_by_manual_ids(self.session,[row.id])
return {
"inbox_id": None,
"manual_upload_id": row.id,
"form_data_id": _as_uuid(form_ids.get(str(row.id))),
"job_post_id": row.job_post_id,
"candidate_user_id": row.user_id,
}
async def _resolve_form(self,form_id,requested_job):
form_row=await FormData.get_form_data_by_id(self.session,form_id)
if not form_row:
raise HTTPException(status_code=404,detail="Form applicant not found")
job_id=form_row.assigned_job_post_id or form_row.job_post_id
if job_id is None:
raise HTTPException(status_code=422,detail="Candidate is not assigned to a job")
if requested_job is not None and job_id!=requested_job:
raise HTTPException(status_code=422,detail="job_post_id does not match the application")
if form_row.is_duplicate:
raise HTTPException(status_code=422,detail="Duplicate form applicants cannot receive an offer")
if (form_row.processing_state or "").strip().lower()=="rejected":
raise HTTPException(status_code=422,detail="Rejected form applicants cannot receive an offer")
if not form_row.job_post_id and form_row.assigned_job_post_id:
form_row.job_post_id=form_row.assigned_job_post_id
self.session.add(form_row)
await self.session.commit()
if form_row.manual_upload_candidate_id:
return await self._resolve_manual(
form_row.manual_upload_candidate_id,job_id,require_interview=False,
)
from g_sheet.views import SheetFormData
promoted=await SheetFormData(session=self.session)._promote_to_application(form_row)
if not promoted or not promoted.user_id:
raise HTTPException(status_code=422,detail="Could not promote this form applicant")
return {
"inbox_id": None,
"manual_upload_id": promoted.id,
"form_data_id": form_row.id,
"job_post_id": promoted.job_post_id or job_id,
"candidate_user_id": promoted.user_id,
}
async def _offer_mail_context(self,row,app):
names=await Users.names_by_ids(self.session,[row.candidate_user_id])
name=names.get(str(row.candidate_user_id))
email=None
result=await self.session.execute(
select(Users.email,Users.name).where(Users.id==row.candidate_user_id)
)
pair=result.first()
if pair:
email=(pair[0] or "").strip().lower() or None
name=name or pair[1]
job=await JobPosts.get_job_post_by_id(self.session,row.job_post_id)
title=job.title if job else None
return name or "Candidate",email,title
async def _advance_pipeline(self,app,current_user):
pipeline=Pipeline(session=self.session)
try:
if app.get("inbox_id") is not None:
await pipeline.change_stage(
"OFFER",current_user,inbox_id=app["inbox_id"],
change_reason="Offer sent",
)
elif app.get("manual_upload_id") is not None:
await pipeline.change_stage(
"OFFER",current_user,manual_upload_id=app["manual_upload_id"],
change_reason="Offer sent",
)
except HTTPException as exc:
if exc.status_code!=400:
raise
async def _notify_send_failed(self,user_id,candidate_name,job_title,row):
try:
bits=[candidate_name] if candidate_name else []
if job_title:
bits.append(job_title)
body=" · ".join(bits) if bits else MAIL_FAIL_DETAIL
await Notifications.insert_notification(self.session,{
"user_id": user_id,
"kind": "message",
"title": MAIL_FAIL_DETAIL,
"body": body,
"link_path": f"/offers?offer={row.id}",
"job_post_id": row.job_post_id,
})
except Exception:
pass