258 lines
11 KiB
Python
258 lines
11 KiB
Python
"""FormData + SheetImportRun — spreadsheet mirror and background import runs."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import uuid
|
|
from datetime import datetime, timezone
|
|
|
|
from sqlalchemy import Column, DateTime, Index, delete, func, insert, or_
|
|
from sqlalchemy.dialects.postgresql import JSONB
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
from sqlmodel import Field, SQLModel, select
|
|
|
|
|
|
def _now() -> datetime:
|
|
return datetime.now(timezone.utc)
|
|
|
|
|
|
_BULK_CHUNK = 1000
|
|
|
|
|
|
class FormData(SQLModel, table=True):
|
|
"""One spreadsheet data row. raw_record keeps the full original header→value map."""
|
|
|
|
__tablename__ = "form_data"
|
|
__table_args__ = (
|
|
Index("ix_form_data_sheet_row_number", "sheet", "row_number", unique=True),
|
|
)
|
|
|
|
id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True)
|
|
sheet: str = Field(nullable=False, index=True)
|
|
# Optional link to a job post. DB FK only — no ORM Relationship (avoids
|
|
# pulling job_posts into the sheet worker metadata graph).
|
|
job_post_id: uuid.UUID | None = Field(default=None, index=True)
|
|
row_number: int | None = Field(default=None)
|
|
serial_no: str | None = Field(default=None)
|
|
entry_year: str | None = Field(default=None)
|
|
entry_month: str | None = Field(default=None)
|
|
entry_date: datetime | None = Field(default=None, sa_type=DateTime(timezone=True))
|
|
entry_time: str | None = Field(default=None)
|
|
screened_by: str | None = Field(default=None, index=True)
|
|
name: str | None = Field(default=None, index=True)
|
|
gender: str | None = Field(default=None)
|
|
date_of_birth: datetime | None = Field(default=None, sa_type=DateTime(timezone=True))
|
|
cnic: str | None = Field(default=None, index=True)
|
|
cgpa: str | None = Field(default=None)
|
|
hr_comments: str | None = Field(default=None)
|
|
candidate_number: str | None = Field(default=None)
|
|
candidate_email: str | None = Field(default=None, index=True)
|
|
profile_link: str | None = Field(default=None)
|
|
resume_link: str | None = Field(default=None)
|
|
area_of_expertise: str | None = Field(default=None)
|
|
requisition_number: str | None = Field(default=None, index=True)
|
|
position_suitable_for: str | None = Field(default=None)
|
|
source_of_application: str | None = Field(default=None)
|
|
age: int | None = Field(default=None)
|
|
age_raw: str | None = Field(default=None)
|
|
marital_status: str | None = Field(default=None)
|
|
degree: str | None = Field(default=None)
|
|
university: str | None = Field(default=None)
|
|
university_other: str | None = Field(default=None)
|
|
experience: str | None = Field(default=None)
|
|
experience_details: str | None = Field(default=None)
|
|
area_of_residence: str | None = Field(default=None)
|
|
residing_city: str | None = Field(default=None)
|
|
residing_country: str | None = Field(default=None)
|
|
communication_skills: int | None = Field(default=None)
|
|
preferred_timings: str | None = Field(default=None)
|
|
ho_availability: str | None = Field(default=None)
|
|
current_company: str | None = Field(default=None)
|
|
reason_for_leaving: str | None = Field(default=None)
|
|
notice_period: str | None = Field(default=None)
|
|
current_salary: str | None = Field(default=None)
|
|
current_salary_value: int | None = Field(default=None)
|
|
expected_salary: str | None = Field(default=None)
|
|
expected_salary_value: int | None = Field(default=None)
|
|
director_poc_category: str | None = Field(default=None)
|
|
pros: str | None = Field(default=None)
|
|
cons: str | None = Field(default=None)
|
|
|
|
raw_record: dict | None = Field(default=None, sa_column=Column(JSONB))
|
|
imported_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True))
|
|
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
|
|
def _filters(cls, *, sheet=None, search=None):
|
|
filters = []
|
|
if sheet:
|
|
filters.append(cls.sheet == sheet)
|
|
if search:
|
|
# Twelve unanchored ILIKEs over ~26k rows is a sequential scan of a few
|
|
# tens of ms — acceptable at this size; a pg_trgm GIN index is the
|
|
# upgrade if the sheet grows an order of magnitude.
|
|
pattern = f"%{search}%"
|
|
filters.append(or_(
|
|
cls.name.ilike(pattern),
|
|
cls.candidate_email.ilike(pattern),
|
|
cls.candidate_number.ilike(pattern),
|
|
cls.screened_by.ilike(pattern),
|
|
cls.degree.ilike(pattern),
|
|
cls.university.ilike(pattern),
|
|
cls.experience.ilike(pattern),
|
|
cls.experience_details.ilike(pattern),
|
|
cls.current_company.ilike(pattern),
|
|
cls.position_suitable_for.ilike(pattern),
|
|
cls.area_of_expertise.ilike(pattern),
|
|
cls.source_of_application.ilike(pattern),
|
|
cls.cnic.ilike(pattern),
|
|
cls.residing_city.ilike(pattern),
|
|
))
|
|
return filters
|
|
|
|
@classmethod
|
|
async def get_form_data_by_id(cls, session: AsyncSession, record_id):
|
|
try:
|
|
rid = uuid.UUID(str(record_id))
|
|
except (TypeError, ValueError):
|
|
return None
|
|
result = await session.execute(select(cls).where(cls.id == rid))
|
|
return result.scalars().first()
|
|
|
|
@classmethod
|
|
async def fetch_form_data(cls, session: AsyncSession, *, sheet=None, search=None, top=None, skip=None):
|
|
statement = select(cls).order_by(cls.sheet, cls.row_number)
|
|
for clause in cls._filters(sheet=sheet, search=search):
|
|
statement = statement.where(clause)
|
|
if skip:
|
|
statement = statement.offset(skip)
|
|
if top is not None:
|
|
statement = statement.limit(top)
|
|
result = await session.execute(statement)
|
|
return result.scalars().all()
|
|
|
|
@classmethod
|
|
async def count_form_data(cls, session: AsyncSession, *, sheet=None, search=None):
|
|
statement = select(func.count()).select_from(cls)
|
|
for clause in cls._filters(sheet=sheet, search=search):
|
|
statement = statement.where(clause)
|
|
result = await session.execute(statement)
|
|
return result.scalar_one()
|
|
|
|
@classmethod
|
|
async def get_sheet_names(cls, session: AsyncSession):
|
|
result = await session.execute(
|
|
select(cls.sheet).distinct().order_by(cls.sheet)
|
|
)
|
|
return list(result.scalars().all())
|
|
|
|
@classmethod
|
|
async def delete_by_sheet(cls, session: AsyncSession, sheet: str, *, commit: bool = True):
|
|
count_result = await session.execute(
|
|
select(func.count()).select_from(cls).where(cls.sheet == sheet)
|
|
)
|
|
deleted = count_result.scalar_one()
|
|
await session.execute(delete(cls).where(cls.sheet == sheet))
|
|
if commit:
|
|
await session.commit()
|
|
return deleted
|
|
|
|
@classmethod
|
|
async def insert_form_data_bulk(
|
|
cls, session: AsyncSession, records: list[dict], *, commit: bool = True,
|
|
):
|
|
# Core insertmanyvalues — building ~26k ORM instances is the slow path.
|
|
# default_factory does not run on Core insert, so stamp timestamps here.
|
|
now = _now()
|
|
total = 0
|
|
for start in range(0, len(records), _BULK_CHUNK):
|
|
chunk = []
|
|
for fields in records[start:start + _BULK_CHUNK]:
|
|
row = dict(fields)
|
|
row.setdefault("id", uuid.uuid4())
|
|
row.setdefault("imported_at", now)
|
|
row.setdefault("created_at", now)
|
|
row.setdefault("updated_at", now)
|
|
chunk.append(row)
|
|
if chunk:
|
|
await session.execute(insert(cls), chunk)
|
|
total += len(chunk)
|
|
if commit:
|
|
await session.commit()
|
|
return total
|
|
|
|
@classmethod
|
|
async def replace_sheet(cls, session: AsyncSession, sheet: str, records: list[dict]):
|
|
"""Delete + insert in one transaction so a mid-insert failure keeps prior rows."""
|
|
deleted = await cls.delete_by_sheet(session, sheet, commit=False)
|
|
inserted = await cls.insert_form_data_bulk(session, records, commit=False)
|
|
await session.commit()
|
|
return {"deleted": deleted, "inserted": inserted}
|
|
|
|
|
|
class SheetImportRun(SQLModel, table=True):
|
|
"""One Google Sheet → FormData import job (Taskiq). Survives tab close."""
|
|
|
|
__tablename__ = "sheet_import_runs"
|
|
|
|
id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True)
|
|
status: str = Field(default="queued", index=True) # queued|running|completed|failed
|
|
task_id: str | None = Field(default=None)
|
|
# Plain UUID — no ORM FK. Importing users.models pulls Users→Inbox relationships
|
|
# that the sheet worker does not load; the DB constraint still enforces integrity.
|
|
created_by: uuid.UUID | None = Field(default=None)
|
|
tab: str | None = Field(default=None) # None = import all tabs
|
|
report: dict | None = Field(default=None, sa_column=Column(JSONB))
|
|
error: str | None = Field(default=None)
|
|
created_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True))
|
|
started_at: datetime | None = Field(default=None, sa_type=DateTime(timezone=True))
|
|
finished_at: datetime | None = Field(default=None, sa_type=DateTime(timezone=True))
|
|
|
|
@staticmethod
|
|
def _as_uuid(record_id) -> uuid.UUID | None:
|
|
if record_id in (None, ""):
|
|
return None
|
|
try:
|
|
return uuid.UUID(str(record_id))
|
|
except ValueError:
|
|
return None
|
|
|
|
@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()
|
|
|
|
@classmethod
|
|
async def get_active(cls, session: AsyncSession):
|
|
result = await session.execute(
|
|
select(cls)
|
|
.where(cls.status.in_(("queued", "running")))
|
|
.order_by(cls.created_at.desc())
|
|
)
|
|
return result.scalars().first()
|
|
|
|
@classmethod
|
|
async def insert_run(cls, session: AsyncSession, fields: dict, *, commit: bool = True):
|
|
row = cls(**fields)
|
|
session.add(row)
|
|
if commit:
|
|
await session.commit()
|
|
await session.refresh(row)
|
|
return row
|
|
|
|
@classmethod
|
|
async def update_run(cls, session: AsyncSession, record_id, fields: dict, *, commit: bool = True):
|
|
row = await cls.get_by_id(session, record_id)
|
|
if not row:
|
|
return None
|
|
for key, value in fields.items():
|
|
setattr(row, key, value)
|
|
session.add(row)
|
|
if commit:
|
|
await session.commit()
|
|
await session.refresh(row)
|
|
return row
|