332 lines
11 KiB
Python
332 lines
11 KiB
Python
"""Drive CV extract wrapper for sheet ingest.
|
|
|
|
Hang `@extract_drive_cvs` on `SheetImport.import_sheet` only (the worker).
|
|
HTTP enqueue routes must not run this — FormData rows do not exist yet.
|
|
|
|
Worker job pattern (same as a Taskiq message): create a temp dir for the run,
|
|
stream each Drive CV to a file, extract, write extracted_data, delete that file.
|
|
A finally block removes the job dir so a successful run leaves no CVs on disk.
|
|
One row at a time — a plain sequential loop, no extra locks.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import logging
|
|
import os
|
|
import shutil
|
|
import uuid
|
|
from datetime import datetime,timezone
|
|
from functools import wraps
|
|
from inspect import signature
|
|
from pathlib import Path
|
|
|
|
from dotenv import load_dotenv
|
|
|
|
from g_sheet.models import FormData
|
|
from g_sheet.plugins import (
|
|
SheetsApiError,
|
|
drive_file_id,
|
|
download_drive_file,
|
|
ensure_fresh,
|
|
load_credentials,
|
|
)
|
|
|
|
load_dotenv(Path(__file__).resolve().parent.parent/".env")
|
|
|
|
logger=logging.getLogger("g_sheet.decorators")
|
|
|
|
_MAX_RESUME_CHARS=int(os.getenv("MAX_RESUME_CHARS","60000"))
|
|
_MAX_PDF_SIZE_MB=int(os.getenv("MAX_PDF_SIZE_MB","10"))
|
|
_TEMP_ROOT=Path(__file__).resolve().parent/"tmp"/"cv_extract"
|
|
|
|
|
|
def build_extracted_data(
|
|
*,
|
|
status,
|
|
resume_link,
|
|
file_id=None,
|
|
filename=None,
|
|
mime_type=None,
|
|
text=None,
|
|
page_count=None,
|
|
truncated=None,
|
|
error_code=None,
|
|
error_message=None,
|
|
):
|
|
"""Stable JSON blob stored on form_data.extracted_data."""
|
|
return {
|
|
"status":status,
|
|
"resume_link":resume_link or "",
|
|
"file_id":file_id,
|
|
"filename":filename,
|
|
"mime_type":mime_type,
|
|
"text":text,
|
|
"page_count":page_count,
|
|
"truncated":truncated,
|
|
"char_count":len(text) if isinstance(text,str) else None,
|
|
"error_code":error_code,
|
|
"error_message":error_message,
|
|
"extracted_at":datetime.now(timezone.utc).isoformat(),
|
|
}
|
|
|
|
|
|
def _drive_error_code(status_code):
|
|
if status_code in (401,403):
|
|
return "DRIVE_FORBIDDEN"
|
|
if status_code==404:
|
|
return "DRIVE_FILE_NOT_FOUND"
|
|
if status_code==413:
|
|
return "PAYLOAD_TOO_LARGE"
|
|
if status_code in (400,415):
|
|
return "UNSUPPORTED_FILE_TYPE"
|
|
return "DRIVE_DOWNLOAD_FAILED"
|
|
|
|
|
|
def _prepare_drive_credentials(service):
|
|
"""Load/refresh the Google session. Never raises — None means skip extract."""
|
|
try:
|
|
if service is None:
|
|
return load_credentials()
|
|
creds=getattr(service,"credentials",None)
|
|
path=getattr(service,"credentials_path",None)
|
|
scopes=getattr(service,"scopes",None)
|
|
if creds is not None:
|
|
return ensure_fresh(creds,path)
|
|
return load_credentials(path,scopes)
|
|
except Exception:
|
|
logger.warning("Google Drive session unavailable; skipping CV extract")
|
|
return None
|
|
|
|
|
|
def _is_sheet_service(obj):
|
|
return obj is not None and hasattr(obj,"session") and hasattr(obj,"spreadsheet_id")
|
|
|
|
|
|
def _tab_from(result,args,kwargs):
|
|
if isinstance(result,dict) and result.get("tab"):
|
|
return result.get("tab")
|
|
if kwargs.get("tab"):
|
|
return kwargs.get("tab")
|
|
if args:
|
|
return args[0]
|
|
return None
|
|
|
|
|
|
def _should_ingest(result):
|
|
if not isinstance(result,dict):
|
|
return False
|
|
if result.get("error"):
|
|
return False
|
|
if result.get("status") in ("queued","running","failed"):
|
|
return False
|
|
if result.get("rows_read",1)==0 and result.get("inserted",1)==0:
|
|
return False
|
|
return True
|
|
|
|
|
|
def make_job_temp_dir(root=None):
|
|
"""Temp dir for one extract job. Caller must remove_job_temp_dir in finally."""
|
|
base=Path(root) if root else _TEMP_ROOT
|
|
base.mkdir(parents=True,exist_ok=True)
|
|
job_dir=base/uuid.uuid4().hex
|
|
job_dir.mkdir()
|
|
return job_dir
|
|
|
|
|
|
def remove_job_temp_dir(job_dir):
|
|
"""Delete leftover CVs and the job dir. No-op if missing."""
|
|
if not job_dir:
|
|
return
|
|
path=Path(job_dir)
|
|
if not path.exists():
|
|
return
|
|
shutil.rmtree(path,ignore_errors=True)
|
|
|
|
|
|
def _unlink(path):
|
|
if path is None:
|
|
return
|
|
try:
|
|
Path(path).unlink(missing_ok=True)
|
|
except OSError as e:
|
|
logger.warning("could not delete temp CV %s: %s",path,e)
|
|
|
|
|
|
def _extract_pdf(data,filename,max_chars):
|
|
from app.services.pdf import extract_resume,sanitize_filename
|
|
from job.candidate.plugins import normalize_spaced_text
|
|
|
|
resume=extract_resume(data,sanitize_filename(filename),max_chars)
|
|
return {
|
|
"text":normalize_spaced_text(resume.text),
|
|
"page_count":resume.page_count,
|
|
"truncated":resume.truncated,
|
|
}
|
|
|
|
|
|
def _download_and_extract(credentials,link,max_chars,max_bytes,dest_dir):
|
|
"""Stream one Drive file into dest_dir, extract, then delete that file."""
|
|
file_id=drive_file_id(link)
|
|
if not file_id:
|
|
return build_extracted_data(
|
|
status="skipped",
|
|
resume_link=link,
|
|
error_code="NOT_DRIVE_URL",
|
|
error_message="Resume link is not a Google Drive file URL",
|
|
)
|
|
dest=None
|
|
try:
|
|
downloaded=download_drive_file(
|
|
credentials,link,max_bytes=max_bytes,dest_dir=dest_dir,
|
|
)
|
|
dest=downloaded.get("path")
|
|
if dest is None:
|
|
return build_extracted_data(
|
|
status="failed",
|
|
resume_link=link,
|
|
file_id=downloaded.get("file_id") or file_id,
|
|
error_code="DRIVE_DOWNLOAD_FAILED",
|
|
error_message="Drive download did not write a file",
|
|
)
|
|
data=Path(dest).read_bytes()
|
|
try:
|
|
parsed=_extract_pdf(data,downloaded.get("filename") or "resume.pdf",max_chars)
|
|
finally:
|
|
data=b""
|
|
return build_extracted_data(
|
|
status="completed",
|
|
resume_link=link,
|
|
file_id=downloaded.get("file_id") or file_id,
|
|
filename=downloaded.get("filename"),
|
|
mime_type=downloaded.get("mime_type"),
|
|
text=parsed["text"],
|
|
page_count=parsed["page_count"],
|
|
truncated=parsed["truncated"],
|
|
)
|
|
except SheetsApiError as e:
|
|
return build_extracted_data(
|
|
status="failed",
|
|
resume_link=link,
|
|
file_id=file_id,
|
|
error_code=_drive_error_code(e.status_code),
|
|
error_message=(e.message or "")[:300],
|
|
)
|
|
except Exception as e:
|
|
from app.core.errors import ATSError
|
|
if isinstance(e,ATSError):
|
|
return build_extracted_data(
|
|
status="failed",
|
|
resume_link=link,
|
|
file_id=file_id,
|
|
error_code=e.error_code,
|
|
error_message=e.public_message,
|
|
)
|
|
logger.exception("drive download/extract failed for file_id=%s",file_id)
|
|
return build_extracted_data(
|
|
status="failed",
|
|
resume_link=link,
|
|
file_id=file_id,
|
|
error_code="DRIVE_DOWNLOAD_FAILED",
|
|
error_message="Drive download failed",
|
|
)
|
|
finally:
|
|
_unlink(dest)
|
|
|
|
|
|
async def extract_one_resume(credentials,resume_link,dest_dir,max_chars=None,max_bytes=None):
|
|
"""Download one Drive URL into dest_dir and return extracted_data JSON."""
|
|
if max_chars is None:
|
|
max_chars=_MAX_RESUME_CHARS
|
|
if max_bytes is None:
|
|
max_bytes=_MAX_PDF_SIZE_MB*1024*1024
|
|
link=(resume_link or "").strip()
|
|
return await asyncio.to_thread(
|
|
_download_and_extract,credentials,link,max_chars,max_bytes,dest_dir,
|
|
)
|
|
|
|
|
|
async def ingest_form_resume_links(session,sheet,credentials,temp_root=None):
|
|
"""One Drive file per resume_link: download → extract → DB → delete file.
|
|
|
|
The job temp dir is created at start and removed in finally so a finished
|
|
run leaves no CVs on disk (Taskiq worker cleanup).
|
|
"""
|
|
rows=await FormData.fetch_resume_links(session,sheet)
|
|
completed=0
|
|
failed=0
|
|
max_chars=_MAX_RESUME_CHARS
|
|
max_bytes=_MAX_PDF_SIZE_MB*1024*1024
|
|
job_dir=make_job_temp_dir(temp_root)
|
|
logger.info(
|
|
"drive CV extract starting tab=%s resumes=%s temp=%s",
|
|
sheet,len(rows),job_dir,
|
|
)
|
|
try:
|
|
for record_id,resume_link in rows:
|
|
payload=await extract_one_resume(
|
|
credentials,resume_link,job_dir,max_chars,max_bytes,
|
|
)
|
|
try:
|
|
await FormData.set_extracted_data(session,record_id,payload)
|
|
except Exception:
|
|
logger.exception("could not persist extracted_data for %s",record_id)
|
|
failed+=1
|
|
try:
|
|
await session.rollback()
|
|
except Exception:
|
|
logger.exception("rollback after extracted_data persist failed")
|
|
continue
|
|
status=payload.get("status")
|
|
if status=="completed":
|
|
completed+=1
|
|
elif status=="failed":
|
|
failed+=1
|
|
logger.info(
|
|
"drive CV extract finished tab=%s extracted=%s failed=%s",
|
|
sheet,completed,failed,
|
|
)
|
|
return {"extracted":completed,"extract_failed":failed}
|
|
finally:
|
|
remove_job_temp_dir(job_dir)
|
|
|
|
|
|
def extract_drive_cvs(func):
|
|
"""Hang on SheetImport.import_sheet (worker ingest), not on HTTP enqueue.
|
|
|
|
After rows are inserted: one Drive download + extract per resume_link,
|
|
written to form_data.extracted_data at the end of each row.
|
|
Credentials load only if ingest will actually run.
|
|
"""
|
|
|
|
@wraps(func)
|
|
async def wrapper(*args,**kwargs):
|
|
result=await func(*args,**kwargs)
|
|
service=args[0] if args and _is_sheet_service(args[0]) else None
|
|
session=getattr(service,"session",None) if service is not None else None
|
|
rest=args[1:] if service is not None else args
|
|
tab=_tab_from(result,rest,kwargs)
|
|
if session is None or not tab or not _should_ingest(result):
|
|
logger.info(
|
|
"drive CV extract skipped tab=%s session=%s ingest=%s",
|
|
tab,session is not None,
|
|
_should_ingest(result) if isinstance(result,dict) else False,
|
|
)
|
|
return result
|
|
credentials=await asyncio.to_thread(_prepare_drive_credentials,service)
|
|
if credentials is None:
|
|
logger.warning("Google Drive session unavailable; skipping CV extract")
|
|
return result
|
|
try:
|
|
stats=await ingest_form_resume_links(session,tab,credentials)
|
|
except Exception:
|
|
logger.exception("drive CV extract after import of %s failed",tab)
|
|
stats={"extracted":0,"extract_failed":0}
|
|
if isinstance(result,dict):
|
|
result["extracted"]=stats.get("extracted",0)
|
|
result["extract_failed"]=stats.get("extract_failed",0)
|
|
return result
|
|
|
|
wrapper.__signature__=signature(func)
|
|
return wrapper
|