639 lines
27 KiB
Python
639 lines
27 KiB
Python
"""Google Sheets service — business logic for the g_sheet domain.
|
|
|
|
The Google client is blocking, so every call goes through asyncio.to_thread rather
|
|
than stalling the event loop. Client construction is lazy and guarded by a lock so
|
|
concurrent requests build it exactly once.
|
|
|
|
Hierarchy:
|
|
Sheet shared config / session
|
|
└─ SheetClient credentials + spreadsheets client
|
|
├─ SheetRead
|
|
│ ├─ SheetHealth
|
|
│ └─ SheetImport
|
|
└─ SheetWrite
|
|
SheetFormData DB mirror only (no Google client)
|
|
"""
|
|
|
|
import asyncio
|
|
import logging
|
|
import threading
|
|
import uuid
|
|
from datetime import datetime,timezone
|
|
|
|
from fastapi import HTTPException
|
|
|
|
from g_sheet.decorators import extract_drive_cvs
|
|
from g_sheet.plugins import (
|
|
SCOPES,
|
|
SPREADSHEET_ID,
|
|
SPREADSHEET_NAME,
|
|
SPREADSHEET_URL,
|
|
SheetsServiceError,
|
|
build_sheets_client,
|
|
ensure_fresh,
|
|
execute,
|
|
load_credentials,
|
|
quote_tab,
|
|
rows_to_indexed_records,
|
|
rows_to_records,
|
|
stringify_rows,
|
|
)
|
|
from g_sheet.models import FormData,SheetImportRun
|
|
from g_sheet.serializers import (
|
|
serialize_append,
|
|
serialize_clear,
|
|
serialize_form_data,
|
|
serialize_health,
|
|
serialize_import,
|
|
serialize_import_all,
|
|
serialize_import_run,
|
|
serialize_metadata,
|
|
serialize_records,
|
|
serialize_sheet_summary,
|
|
serialize_update,
|
|
serialize_values,
|
|
)
|
|
|
|
logger=logging.getLogger("g_sheet.views")
|
|
|
|
|
|
class Sheet:
|
|
"""Parent: spreadsheet identity, optional DB session, and shared helpers."""
|
|
|
|
def __init__(self,session=None,spreadsheet_id=None,credentials_path=None,scopes=None):
|
|
self.session=session
|
|
self.spreadsheet_id=spreadsheet_id or SPREADSHEET_ID
|
|
self.spreadsheet_name=SPREADSHEET_NAME
|
|
self.spreadsheet_url=SPREADSHEET_URL
|
|
self.credentials_path=credentials_path
|
|
self.scopes=scopes or SCOPES
|
|
self.credentials=None
|
|
self.client=None
|
|
self._lock=threading.Lock()
|
|
|
|
def _require_session(self):
|
|
if self.session is None:
|
|
raise HTTPException(status_code=500,detail="Database session is required")
|
|
return self.session
|
|
|
|
|
|
class SheetClient(Sheet):
|
|
"""Google API client — lazy connect, token refresh, values/spreadsheets handles."""
|
|
|
|
def _connect(self):
|
|
"""Build credentials + client once, then keep refreshing the same token.
|
|
|
|
Double-checked under the lock: two requests racing here must not each build
|
|
their own client.
|
|
"""
|
|
if self.client is not None:
|
|
return ensure_fresh(self.credentials,self.credentials_path) and self.client
|
|
with self._lock:
|
|
if self.client is None:
|
|
self.credentials=load_credentials(self.credentials_path,self.scopes)
|
|
self.client=build_sheets_client(self.credentials)
|
|
else:
|
|
ensure_fresh(self.credentials,self.credentials_path)
|
|
return self.client
|
|
|
|
async def _values(self):
|
|
if not self.spreadsheet_id:
|
|
raise HTTPException(status_code=500,detail="SPREADSHEET_ID is not configured")
|
|
client=await asyncio.to_thread(self._connect)
|
|
return client.spreadsheets().values()
|
|
|
|
async def _spreadsheets(self):
|
|
if not self.spreadsheet_id:
|
|
raise HTTPException(status_code=500,detail="SPREADSHEET_ID is not configured")
|
|
client=await asyncio.to_thread(self._connect)
|
|
return client.spreadsheets()
|
|
|
|
|
|
class SheetRead(SheetClient):
|
|
"""Read-only sheet operations."""
|
|
|
|
async def get_metadata(self):
|
|
"""Spreadsheet title, id, url and every tab with its row/column counts."""
|
|
try:
|
|
spreadsheets=await self._spreadsheets()
|
|
request=spreadsheets.get(spreadsheetId=self.spreadsheet_id,fields=(
|
|
"spreadsheetId,spreadsheetUrl,properties(title,locale,timeZone),"
|
|
"sheets(properties(sheetId,title,index,gridProperties(rowCount,columnCount)))"
|
|
))
|
|
payload=await asyncio.to_thread(execute,request,"spreadsheet metadata")
|
|
return serialize_metadata(payload)
|
|
except SheetsServiceError as e:
|
|
raise HTTPException(status_code=e.status_code,detail=e.message)
|
|
|
|
async def list_tabs(self):
|
|
"""Tab titles in sheet order."""
|
|
metadata=await self.get_metadata()
|
|
return [tab["title"] for tab in metadata["tabs"] if tab.get("title")]
|
|
|
|
async def read_range(self,tab,cell_range=None):
|
|
"""Raw rows for a tab, or for a sub-range of it when cell_range is given."""
|
|
try:
|
|
values=await self._values()
|
|
target=quote_tab(tab,cell_range)
|
|
request=values.get(spreadsheetId=self.spreadsheet_id,range=target)
|
|
payload=await asyncio.to_thread(execute,request,f"read {target}")
|
|
rows=stringify_rows(payload.get("values"))
|
|
return serialize_values(tab,cell_range,rows)
|
|
except SheetsServiceError as e:
|
|
raise HTTPException(status_code=e.status_code,detail=e.message)
|
|
|
|
async def read_records(self,tab):
|
|
"""Rows keyed by the first row. Blank rows are skipped, short rows padded."""
|
|
data=await self.read_range(tab)
|
|
return serialize_records(tab,rows_to_records(data["rows"]))
|
|
|
|
async def read_all(self):
|
|
"""Every tab as records, keyed by tab name."""
|
|
tabs=await self.list_tabs()
|
|
sheets={}
|
|
for tab in tabs:
|
|
data=await self.read_records(tab)
|
|
sheets[tab]=data["records"]
|
|
return {"sheets":sheets,"tabs":tabs,"total":len(tabs)}
|
|
|
|
|
|
class SheetWrite(SheetClient):
|
|
"""Mutating sheet operations."""
|
|
|
|
async def append_rows(self,tab,rows):
|
|
"""Append rows below the tab's current content."""
|
|
if not rows:
|
|
raise HTTPException(status_code=422,detail="rows must not be empty")
|
|
try:
|
|
values=await self._values()
|
|
target=quote_tab(tab)
|
|
request=values.append(
|
|
spreadsheetId=self.spreadsheet_id,
|
|
range=target,
|
|
valueInputOption="USER_ENTERED",
|
|
insertDataOption="INSERT_ROWS",
|
|
body={"values":rows},
|
|
)
|
|
payload=await asyncio.to_thread(execute,request,f"append to {target}")
|
|
return serialize_append(tab,payload)
|
|
except SheetsServiceError as e:
|
|
raise HTTPException(status_code=e.status_code,detail=e.message)
|
|
|
|
async def update_range(self,tab,cell_range,rows):
|
|
"""Overwrite an explicit A1 range with rows."""
|
|
if not cell_range:
|
|
raise HTTPException(status_code=422,detail="cell_range is required")
|
|
if not rows:
|
|
raise HTTPException(status_code=422,detail="rows must not be empty")
|
|
try:
|
|
values=await self._values()
|
|
target=quote_tab(tab,cell_range)
|
|
request=values.update(
|
|
spreadsheetId=self.spreadsheet_id,
|
|
range=target,
|
|
valueInputOption="USER_ENTERED",
|
|
body={"values":rows},
|
|
)
|
|
payload=await asyncio.to_thread(execute,request,f"update {target}")
|
|
return serialize_update(tab,payload)
|
|
except SheetsServiceError as e:
|
|
raise HTTPException(status_code=e.status_code,detail=e.message)
|
|
|
|
async def clear_range(self,tab,cell_range):
|
|
"""Clear the values in an explicit A1 range, leaving formatting intact."""
|
|
if not cell_range:
|
|
raise HTTPException(status_code=422,detail="cell_range is required")
|
|
try:
|
|
values=await self._values()
|
|
target=quote_tab(tab,cell_range)
|
|
request=values.clear(spreadsheetId=self.spreadsheet_id,range=target,body={})
|
|
payload=await asyncio.to_thread(execute,request,f"clear {target}")
|
|
return serialize_clear(tab,payload)
|
|
except SheetsServiceError as e:
|
|
raise HTTPException(status_code=e.status_code,detail=e.message)
|
|
|
|
|
|
class SheetHealth(SheetRead):
|
|
"""Credentials + spreadsheet reachability."""
|
|
|
|
async def health_check(self):
|
|
"""Credentials + sheet reachability as a status dict. Never raises."""
|
|
if not self.spreadsheet_id:
|
|
return serialize_health(False,"SPREADSHEET_ID is not configured")
|
|
try:
|
|
tabs=await self.list_tabs()
|
|
return serialize_health(True,"spreadsheet reachable",tabs)
|
|
except HTTPException as e:
|
|
logger.warning("sheets health check failed: %s",e.detail)
|
|
return serialize_health(False,str(e.detail))
|
|
except Exception as e:
|
|
logger.warning("sheets health check failed: %s",e)
|
|
return serialize_health(False,str(e))
|
|
|
|
|
|
class SheetImport(SheetRead):
|
|
"""Google Sheet → FormData import + import-run tracking."""
|
|
|
|
@extract_drive_cvs
|
|
async def import_sheet(self,tab):
|
|
"""Read one tab from Google Sheets and replace its FormData rows."""
|
|
session=self._require_session()
|
|
if not tab or not str(tab).strip():
|
|
raise HTTPException(status_code=422,detail="tab is required")
|
|
tab=str(tab).strip()
|
|
data=await self.read_range(tab)
|
|
rows=data["rows"]
|
|
if not rows:
|
|
return serialize_import({"tab":tab,"rows_read":0,"inserted":0,"deleted":0})
|
|
indexed=rows_to_indexed_records(rows)
|
|
mapped=[
|
|
FormData.from_sheet_row(tab,row_number,record)
|
|
for row_number,record in indexed
|
|
]
|
|
mapped=await FormData.stamp_suggested_job_posts(session,mapped)
|
|
result=await FormData.replace_sheet(session,tab,mapped)
|
|
from inbox.views import Reapplied
|
|
await Reapplied(session=session).sync_for_emails(
|
|
[r.get("candidate_email") for r in mapped]
|
|
)
|
|
return serialize_import({
|
|
"tab":tab,
|
|
"rows_read":len(indexed),
|
|
"inserted":result["inserted"],
|
|
"deleted":result["deleted"],
|
|
})
|
|
|
|
async def import_all(self):
|
|
"""Import every tab sequentially; one tab failure does not abort the rest."""
|
|
self._require_session()
|
|
tabs=await self.list_tabs()
|
|
reports=[]
|
|
for tab in tabs:
|
|
try:
|
|
report=await self.import_sheet(tab)
|
|
reports.append(report)
|
|
except HTTPException as e:
|
|
logger.warning("import_all tab %s failed: %s",tab,e.detail)
|
|
reports.append(serialize_import({
|
|
"tab":tab,"rows_read":0,"inserted":0,"deleted":0,
|
|
"error":str(e.detail),
|
|
}))
|
|
except Exception as e:
|
|
logger.exception("import_all tab %s failed",tab)
|
|
reports.append(serialize_import({
|
|
"tab":tab,"rows_read":0,"inserted":0,"deleted":0,
|
|
"error":str(e),
|
|
}))
|
|
return serialize_import_all(reports)
|
|
|
|
async def start_import(self,current_user=None,tab=None):
|
|
"""Enqueue a sheet import.
|
|
|
|
At the start of every new job: queued/running → keep that job;
|
|
failed → delete those rows and start this one; completed → start this one.
|
|
"""
|
|
session=self._require_session()
|
|
active=await SheetImportRun.get_active(session)
|
|
if active:
|
|
return serialize_import_run(active)
|
|
|
|
await SheetImportRun.delete_failed(session)
|
|
|
|
created_by=None
|
|
if isinstance(current_user,dict) and current_user.get("id"):
|
|
created_by=SheetImportRun._as_uuid(current_user.get("id"))
|
|
|
|
tab_value=str(tab).strip() if tab else None
|
|
row=await SheetImportRun.insert_run(session,{
|
|
"status":"queued",
|
|
"created_by":created_by,
|
|
"tab":tab_value,
|
|
})
|
|
|
|
from g_sheet.tasks import import_sheets
|
|
from taskiq_management.g_sheet_broker_setup import SHEET_QUEUE_NAME
|
|
task=await import_sheets.kicker().with_labels(
|
|
created_at=datetime.now(timezone.utc).isoformat(),
|
|
correlation_id=str(row.id),
|
|
queue=SHEET_QUEUE_NAME,
|
|
).kiq(str(row.id))
|
|
row=await SheetImportRun.update_run(session,row.id,{"task_id":task.task_id})
|
|
return serialize_import_run(row)
|
|
|
|
async def get_import_run(self,run_id=None):
|
|
session=self._require_session()
|
|
if run_id:
|
|
row=await SheetImportRun.get_by_id(session,run_id)
|
|
if not row:
|
|
raise HTTPException(status_code=404,detail="Import run not found")
|
|
return serialize_import_run(row)
|
|
row=await SheetImportRun.get_active(session)
|
|
if row:
|
|
return serialize_import_run(row)
|
|
row=await SheetImportRun.get_latest(session)
|
|
if not row:
|
|
raise HTTPException(status_code=404,detail="No import runs yet")
|
|
return serialize_import_run(row)
|
|
|
|
|
|
class SheetFormData(Sheet):
|
|
"""FormData DB mirror — query / delete only (no Google client)."""
|
|
|
|
async def _hydrate_job_posts(self,items):
|
|
"""Attach suggested job titles, assigned_job_post, and per-job ATS scores.
|
|
|
|
Preferred source is suggested_job_post_ids (ILIKE matches stored on
|
|
import). Legacy rows without that list still title-match. ATS is one
|
|
current score per (form, job). Full JD loads when a card is expanded.
|
|
"""
|
|
if not items:
|
|
return items
|
|
from g_sheet.scoring import serialize_form_ats
|
|
from inbox.models import AtsResults
|
|
from job.job_post.models import JobPosts
|
|
from job.job_post.serializers import serialize_job_post_title
|
|
|
|
session=self._require_session()
|
|
|
|
def _job_payload(post):
|
|
payload=serialize_job_post_title(post)
|
|
if post.is_deleted or not post.is_active:
|
|
payload={**payload,"unavailable":True}
|
|
return payload
|
|
|
|
suggested_ids=[]
|
|
for item in items:
|
|
for raw in item.get("suggested_job_post_ids") or []:
|
|
if raw:
|
|
suggested_ids.append(raw)
|
|
assigned_ids=[]
|
|
for item in items:
|
|
aid=item.get("assigned_job_post_id") or item.get("job_post_id")
|
|
if aid:
|
|
assigned_ids.append(aid)
|
|
wanted=list(dict.fromkeys([*suggested_ids,*assigned_ids]))
|
|
by_id={}
|
|
if wanted:
|
|
for post in await JobPosts.titles_by_ids(session,wanted,active_only=False):
|
|
by_id[str(post.id)]=_job_payload(post)
|
|
|
|
titles=[(item.get("position_applied_for") or "").strip() for item in items]
|
|
titles=[t for t in titles if t]
|
|
by_title={}
|
|
needs_title=any(not (item.get("suggested_job_post_ids") or []) for item in items)
|
|
if titles and needs_title:
|
|
for post in await JobPosts.get_by_titles(session,titles):
|
|
key=(post.title or "").strip().lower()
|
|
by_title.setdefault(key,[]).append(_job_payload(post))
|
|
|
|
ats_by_form=await AtsResults.get_current_for_forms(
|
|
session,[item.get("id") for item in items],
|
|
)
|
|
|
|
for item in items:
|
|
suggested=[str(raw) for raw in (item.get("suggested_job_post_ids") or []) if raw]
|
|
item["suggested_job_post_ids"]=suggested
|
|
if suggested:
|
|
posts=[]
|
|
for sid in suggested:
|
|
payload=by_id.get(sid)
|
|
if payload is None:
|
|
posts.append({"id":sid,"unavailable":True})
|
|
else:
|
|
posts.append(dict(payload))
|
|
item["job_posts"]=posts
|
|
else:
|
|
key=(item.get("position_applied_for") or "").strip().lower()
|
|
item["job_posts"]=[dict(p) for p in (by_title.get(key) or [])]
|
|
|
|
aid=item.get("assigned_job_post_id") or item.get("job_post_id")
|
|
item["assigned_job_post_id"]=str(aid) if aid else None
|
|
item["assigned_job_post"]=by_id.get(str(aid)) if aid else None
|
|
|
|
fid=item.get("id")
|
|
try:
|
|
form_uid=uuid.UUID(str(fid)) if fid else None
|
|
except (TypeError,ValueError):
|
|
form_uid=None
|
|
scores=[serialize_form_ats(row) for row in (ats_by_form.get(form_uid) or [])]
|
|
item["ats_results"]=scores
|
|
score_by_job={
|
|
str(s["job_post_id"]):s for s in scores if s.get("job_post_id")
|
|
}
|
|
for post in item["job_posts"]:
|
|
hit=score_by_job.get(str(post.get("id")))
|
|
if hit:
|
|
post["overall_score"]=hit.get("overall_score")
|
|
post["band"]=hit.get("band")
|
|
assigned_score=score_by_job.get(str(aid)) if aid else None
|
|
if assigned_score and assigned_score.get("overall_score") is not None:
|
|
item["ats_score"]=round(float(assigned_score.get("overall_score")))
|
|
else:
|
|
nums=[s.get("overall_score") for s in scores if s.get("overall_score") is not None]
|
|
item["ats_score"]=round(float(max(nums))) if nums else None
|
|
return items
|
|
|
|
async def get_form_data(
|
|
self,sheet=None,search=None,offset=0,limit=None,
|
|
processing_state=None,is_duplicate=None,has_linkedin=None,has_resume=None,
|
|
city=None,source=None,assigned=None,no_suggestions=None,
|
|
has_suggestions=None,job_post_ids=None,
|
|
):
|
|
session=self._require_session()
|
|
rows=await FormData.fetch_form_data(
|
|
session,sheet=sheet,search=search,offset=offset,limit=limit,
|
|
processing_state=processing_state,is_duplicate=is_duplicate,
|
|
has_linkedin=has_linkedin,has_resume=has_resume,city=city,
|
|
source=source,assigned=assigned,no_suggestions=no_suggestions,
|
|
has_suggestions=has_suggestions,job_post_ids=job_post_ids,
|
|
)
|
|
total=await FormData.count_form_data(
|
|
session,sheet=sheet,search=search,
|
|
processing_state=processing_state,is_duplicate=is_duplicate,
|
|
has_linkedin=has_linkedin,has_resume=has_resume,city=city,
|
|
source=source,assigned=assigned,no_suggestions=no_suggestions,
|
|
has_suggestions=has_suggestions,job_post_ids=job_post_ids,
|
|
)
|
|
items=await self._hydrate_job_posts([serialize_form_data(row) for row in rows])
|
|
from job.candidate.views import CandidateView
|
|
items=await CandidateView(session=session).attach_application_history(items)
|
|
return items,total
|
|
|
|
async def get_form_data_by_id(self,record_id):
|
|
session=self._require_session()
|
|
row=await FormData.get_form_data_by_id(session,record_id)
|
|
if not row:
|
|
raise HTTPException(status_code=404,detail="Form data not found")
|
|
items=await self._hydrate_job_posts([serialize_form_data(row)])
|
|
from job.candidate.views import CandidateView
|
|
return await CandidateView(session=session).attach_application_history(items[0])
|
|
|
|
async def assign_job_post(self,record_id,job_post_id):
|
|
"""Set or clear form_data.assigned_job_post_id (same contract as inbox assign).
|
|
|
|
Setting a job promotes the row into Users + manual_upload_candidate so
|
|
Candidates / Talent Pool / Pipeline can see it (platform tag: Form).
|
|
"""
|
|
session=self._require_session()
|
|
if job_post_id is not None:
|
|
from job.job_post.models import JobPosts
|
|
post=await JobPosts.get_job_post_by_id(session,job_post_id)
|
|
if not post or post.is_deleted or not post.is_active:
|
|
raise HTTPException(status_code=404,detail="Job post not found")
|
|
updated=await FormData.set_job_post(session,record_id,job_post_id)
|
|
if not updated:
|
|
raise HTTPException(status_code=404,detail="Form data not found")
|
|
from inbox.views import Reapplied
|
|
await Reapplied(session=session).sync_for_email(updated.candidate_email)
|
|
if job_post_id is not None:
|
|
await self._promote_to_application(updated)
|
|
from g_sheet.scoring import enqueue_form_score
|
|
await enqueue_form_score(updated.id,job_post_id)
|
|
return await self.get_form_data_by_id(record_id)
|
|
|
|
async def set_processing_state(self,record_id,processing_state,current_user=None):
|
|
allowed=("unread","imported","processed","rejected")
|
|
if processing_state not in allowed:
|
|
raise HTTPException(status_code=422,detail=f"processing_state must be one of {', '.join(allowed)}")
|
|
session=self._require_session()
|
|
row=await FormData.get_form_data_by_id(session,record_id)
|
|
if not row:
|
|
raise HTTPException(status_code=404,detail="Form data not found")
|
|
# Shortlist requires a job — promote (idempotent) then flip the queue label.
|
|
if processing_state=="processed":
|
|
if not row.job_post_id:
|
|
raise HTTPException(status_code=422,detail="Assign a job post before moving to shortlist")
|
|
promoted=await self._promote_to_application(row)
|
|
status=(getattr(promoted,"status",None) or "").strip()
|
|
if promoted is not None and status in ("","CLOSED","PROCESS","BANKED","REJECTED"):
|
|
from job.pipeline.views import Pipeline
|
|
try:
|
|
await Pipeline(session).change_stage(
|
|
"PENDING",current_user,manual_upload_id=promoted.id,
|
|
change_reason="Moved to shortlist from sheet forms",
|
|
)
|
|
except HTTPException as exc:
|
|
if exc.status_code!=400:
|
|
raise
|
|
updated=await FormData.set_processing_state(session,record_id,processing_state)
|
|
if not updated:
|
|
raise HTTPException(status_code=404,detail="Form data not found")
|
|
return await self.get_form_data_by_id(record_id)
|
|
|
|
async def _promote_to_application(self,form_row):
|
|
"""Create Users + manual_upload_candidate from a form_data row (idempotent).
|
|
|
|
Pipeline / Candidates / Talent Pool all read manual_upload_candidate (or
|
|
the CANDIDATE user it creates). platform='Form' is the source badge.
|
|
"""
|
|
session=self._require_session()
|
|
from employment_agent.plugins import parse_linkedin,parse_phone
|
|
from job.candidate.models import Manual_UPLOAD_CANDIDATE
|
|
from job.history.views import HistoryRecorder
|
|
from job.history.enums import HistoryEvent
|
|
|
|
if getattr(form_row,"manual_upload_candidate_id",None):
|
|
existing=await Manual_UPLOAD_CANDIDATE.get_by_id(session,form_row.manual_upload_candidate_id)
|
|
if existing:
|
|
if form_row.job_post_id and existing.job_post_id!=form_row.job_post_id:
|
|
existing.job_post_id=form_row.job_post_id
|
|
session.add(existing)
|
|
await session.commit()
|
|
return existing
|
|
|
|
email=(form_row.candidate_email or "").strip().lower()
|
|
if not email:
|
|
raise HTTPException(status_code=422,detail="candidate_email is required to promote this form applicant")
|
|
if not form_row.job_post_id:
|
|
raise HTTPException(status_code=422,detail="job_post_id is required to promote this form applicant")
|
|
|
|
existing=await Manual_UPLOAD_CANDIDATE.get_by_email_and_job(
|
|
session,email,form_row.job_post_id,
|
|
)
|
|
if existing:
|
|
await FormData.link_manual_upload(session,form_row.id,existing.id)
|
|
return existing
|
|
|
|
resume=(form_row.resume_link or "").strip()
|
|
file_name=""
|
|
if resume:
|
|
file_name=resume.rsplit("/",1)[-1][:180] or "resume"
|
|
|
|
profile=(form_row.profile_link or "").strip()
|
|
linkedin_url=parse_linkedin({"linkedin_url":profile},"").get("linkedin_url")
|
|
phone_fields=parse_phone({"phone":(form_row.candidate_number or "").strip()},"")
|
|
|
|
row=await Manual_UPLOAD_CANDIDATE.create_manual_upload_candidate(session,{
|
|
"candidate_email":email,
|
|
"candidate_name":(form_row.name or "").strip() or email,
|
|
"candidate_phone":phone_fields.get("phone") or "",
|
|
"job_post_id":str(form_row.job_post_id),
|
|
"current_company":(form_row.current_company or "").strip(),
|
|
"current_position":(form_row.position_applied_for or "").strip(),
|
|
"platform":"Form",
|
|
"apply_via":"form",
|
|
"experience":(form_row.experience or "").strip(),
|
|
"status":"PENDING",
|
|
"file_name":file_name,
|
|
"file_path":resume,
|
|
"full_text":"",
|
|
"linkedin_url":linkedin_url,
|
|
})
|
|
from inbox.views import Reapplied
|
|
await Reapplied(session=session).sync_for_email(email)
|
|
await FormData.link_manual_upload(session,form_row.id,row.id)
|
|
try:
|
|
await HistoryRecorder(session).record(
|
|
HistoryEvent.CANDIDATE_CREATED.value,
|
|
actor_id=None,user_id=row.user_id,
|
|
manual_upload_candidate_id=row.id,
|
|
entity_type="manual_upload_candidate",entity_id=row.id,
|
|
to_value=row.candidate_email,
|
|
description="Form",commit=True,
|
|
)
|
|
except Exception:
|
|
logger.exception("form promote history record failed for %s",form_row.id)
|
|
return row
|
|
|
|
async def set_duplicate(self,record_id,is_duplicate):
|
|
if not isinstance(is_duplicate,bool):
|
|
raise HTTPException(status_code=422,detail="is_duplicate must be a boolean")
|
|
updated=await FormData.set_duplicate(self._require_session(),record_id,is_duplicate)
|
|
if not updated:
|
|
raise HTTPException(status_code=404,detail="Form data not found")
|
|
return await self.get_form_data_by_id(record_id)
|
|
|
|
async def set_rating(self,record_id,payload):
|
|
fields={k:v for k,v in (payload or {}).items() if k in ("favorite","rating") and v is not None}
|
|
if not fields:
|
|
raise HTTPException(status_code=400,detail="favorite or rating is required")
|
|
updated=await FormData.set_rating_favorite(self._require_session(),record_id,**fields)
|
|
if not updated:
|
|
raise HTTPException(status_code=404,detail="Form data not found")
|
|
return await self.get_form_data_by_id(record_id)
|
|
|
|
async def get_counts(self,sheet=None,search=None,has_linkedin=None,has_resume=None,city=None,source=None,assigned=None,job_post_ids=None):
|
|
return await FormData.count_processing(
|
|
self._require_session(),sheet=sheet,search=search,
|
|
has_linkedin=has_linkedin,has_resume=has_resume,city=city,
|
|
source=source,assigned=assigned,job_post_ids=job_post_ids,
|
|
)
|
|
|
|
async def count_rows(self,sheet=None,search=None,has_linkedin=None,has_resume=None):
|
|
return await FormData.count_form_data(
|
|
self._require_session(),sheet=sheet,search=search,
|
|
has_linkedin=has_linkedin,has_resume=has_resume,
|
|
)
|
|
|
|
async def get_imported_sheets(self):
|
|
session=self._require_session()
|
|
sheets=await FormData.get_sheet_names(session)
|
|
return serialize_sheet_summary(sheets)
|
|
|
|
async def delete_sheet_data(self,tab):
|
|
session=self._require_session()
|
|
if not tab or not str(tab).strip():
|
|
raise HTTPException(status_code=422,detail="tab is required")
|
|
deleted=await FormData.delete_by_sheet(session,str(tab).strip())
|
|
return {"tab":str(tab).strip(),"deleted":deleted}
|