602 lines
25 KiB
Python
602 lines
25 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, case, 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)
|
|
# Recruiter-assigned job. 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)
|
|
assigned_job_post_id: uuid.UUID | None = Field(default=None, index=True)
|
|
# ILIKE title matches from Position Applied For. One form row → many jobs.
|
|
# Suggested, not assigned. ATS scores each id separately.
|
|
suggested_job_post_ids: list[str] | None = Field(
|
|
default=None, sa_column=Column(JSONB),
|
|
)
|
|
# Set when this form row is promoted into the hiring pipeline (Users +
|
|
# manual_upload_candidate). Idempotency key for assign / shortlist.
|
|
manual_upload_candidate_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)
|
|
# Drive CV extract JSON written by @extract_drive_cvs after sheet ingest.
|
|
extracted_data: dict | None = Field(default=None, sa_column=Column(JSONB))
|
|
area_of_expertise: str | None = Field(default=None)
|
|
requisition_number: str | None = Field(default=None, index=True)
|
|
position_applied_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)
|
|
|
|
# Same allowlist as inbox_messages.processing_state: unread|imported|processed|rejected.
|
|
# server_default is load-bearing — ALTER on a populated form_data table.
|
|
processing_state: str = Field(default="unread", sa_column_kwargs={"server_default": "unread"})
|
|
is_duplicate: bool = Field(default=False, sa_column_kwargs={"server_default": "false"})
|
|
|
|
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, processing_state=None, is_duplicate=None):
|
|
filters = []
|
|
if sheet:
|
|
filters.append(cls.sheet == sheet)
|
|
if processing_state:
|
|
filters.append(cls.processing_state == processing_state)
|
|
if is_duplicate is not None:
|
|
filters.append(cls.is_duplicate == bool(is_duplicate))
|
|
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_applied_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 get_with_job(cls, session: AsyncSession, record_id, job_post_id):
|
|
"""Form row + one job it may be scored against (suggested or assigned)."""
|
|
from job.job_post.models import JobPosts
|
|
|
|
form = await cls.get_form_data_by_id(session, record_id)
|
|
if form is None:
|
|
return None, None
|
|
try:
|
|
jid = uuid.UUID(str(job_post_id))
|
|
except (TypeError, ValueError):
|
|
return None, None
|
|
if jid not in set(cls.score_job_ids(form)):
|
|
return None, None
|
|
job = await JobPosts.get_job_post_by_id(session, jid)
|
|
if job is None or job.is_deleted:
|
|
return None, None
|
|
return form, job
|
|
|
|
@staticmethod
|
|
def score_job_ids(row) -> list[uuid.UUID]:
|
|
"""Jobs ATS may score: assigned only, else every suggested id."""
|
|
if row is None:
|
|
return []
|
|
getter = row.get if isinstance(row, dict) else lambda key, default=None: getattr(row, key, default)
|
|
assigned = getter("assigned_job_post_id") or getter("job_post_id")
|
|
if assigned not in (None, ""):
|
|
try:
|
|
return [uuid.UUID(str(assigned))]
|
|
except (TypeError, ValueError):
|
|
return []
|
|
out: list[uuid.UUID] = []
|
|
seen: set[uuid.UUID] = set()
|
|
for raw in getter("suggested_job_post_ids") or []:
|
|
try:
|
|
uid = uuid.UUID(str(raw))
|
|
except (TypeError, ValueError):
|
|
continue
|
|
if uid not in seen:
|
|
seen.add(uid)
|
|
out.append(uid)
|
|
return out
|
|
|
|
@classmethod
|
|
async def set_job_post(cls, session: AsyncSession, record_id, job_post_id):
|
|
"""Set or clear the recruiter assignment; returns the row or None if missing."""
|
|
row = await cls.get_form_data_by_id(session, record_id)
|
|
if not row:
|
|
return None
|
|
if job_post_id is None:
|
|
row.job_post_id = None
|
|
row.assigned_job_post_id = None
|
|
else:
|
|
try:
|
|
uid = uuid.UUID(str(job_post_id))
|
|
except (TypeError, ValueError):
|
|
return None
|
|
row.job_post_id = uid
|
|
row.assigned_job_post_id = uid
|
|
row.updated_at = _now()
|
|
session.add(row)
|
|
await session.commit()
|
|
await session.refresh(row)
|
|
return row
|
|
|
|
@classmethod
|
|
async def set_processing_state(cls, session: AsyncSession, record_id, processing_state: str):
|
|
row = await cls.get_form_data_by_id(session, record_id)
|
|
if not row:
|
|
return None
|
|
row.processing_state = processing_state
|
|
row.updated_at = _now()
|
|
session.add(row)
|
|
await session.commit()
|
|
await session.refresh(row)
|
|
return row
|
|
|
|
@classmethod
|
|
async def set_duplicate(cls, session: AsyncSession, record_id, is_duplicate: bool):
|
|
row = await cls.get_form_data_by_id(session, record_id)
|
|
if not row:
|
|
return None
|
|
row.is_duplicate = bool(is_duplicate)
|
|
row.updated_at = _now()
|
|
session.add(row)
|
|
await session.commit()
|
|
await session.refresh(row)
|
|
return row
|
|
|
|
@classmethod
|
|
async def link_manual_upload(cls, session: AsyncSession, record_id, manual_upload_candidate_id, *, commit: bool = True):
|
|
row = await cls.get_form_data_by_id(session, record_id)
|
|
if not row:
|
|
return None
|
|
try:
|
|
row.manual_upload_candidate_id = uuid.UUID(str(manual_upload_candidate_id))
|
|
except (TypeError, ValueError):
|
|
return None
|
|
row.updated_at = _now()
|
|
session.add(row)
|
|
if commit:
|
|
await session.commit()
|
|
await session.refresh(row)
|
|
return row
|
|
|
|
@classmethod
|
|
async def set_extracted_data(
|
|
cls, session: AsyncSession, record_id, extracted_data, *, commit: bool = True,
|
|
):
|
|
row = await cls.get_form_data_by_id(session, record_id)
|
|
if not row:
|
|
return None
|
|
row.extracted_data = extracted_data
|
|
row.updated_at = _now()
|
|
session.add(row)
|
|
if commit:
|
|
await session.commit()
|
|
await session.refresh(row)
|
|
return row
|
|
|
|
@classmethod
|
|
async def fetch_resume_links(cls, session: AsyncSession, sheet: str):
|
|
"""(id, resume_link) for one tab. Blank links are dropped."""
|
|
statement = (
|
|
select(cls.id, cls.resume_link)
|
|
.where(cls.sheet == sheet)
|
|
.where(cls.resume_link.is_not(None))
|
|
.order_by(cls.row_number)
|
|
)
|
|
result = await session.execute(statement)
|
|
rows = []
|
|
for record_id, link in result.all():
|
|
text = (link or "").strip()
|
|
if text:
|
|
rows.append((record_id, text))
|
|
return rows
|
|
|
|
@classmethod
|
|
async def fetch_form_data(
|
|
cls, session: AsyncSession, *, sheet=None, search=None,
|
|
processing_state=None, is_duplicate=None, offset=0, limit=None,
|
|
):
|
|
statement = select(cls).order_by(cls.sheet, cls.row_number)
|
|
for clause in cls._filters(
|
|
sheet=sheet, search=search,
|
|
processing_state=processing_state, is_duplicate=is_duplicate,
|
|
):
|
|
statement = statement.where(clause)
|
|
if offset:
|
|
statement = statement.offset(offset)
|
|
if limit is not None:
|
|
statement = statement.limit(limit)
|
|
statement = statement.order_by(cls.row_number)
|
|
result = await session.execute(statement)
|
|
return result.scalars().all()
|
|
|
|
@classmethod
|
|
async def list_by_emails(cls, session: AsyncSession, emails):
|
|
"""Sheet applicants for these addresses. Promoted rows are omitted —
|
|
those already live on manual_upload_candidate."""
|
|
from job.job_post.models import JobPosts
|
|
|
|
lowers = sorted({(e or "").strip().lower() for e in (emails or []) if (e or "").strip()})
|
|
if not lowers:
|
|
return []
|
|
assigned = func.coalesce(cls.assigned_job_post_id, cls.job_post_id)
|
|
result = await session.execute(
|
|
select(cls, JobPosts.title)
|
|
.outerjoin(JobPosts, assigned == JobPosts.id)
|
|
.where(func.lower(cls.candidate_email).in_(lowers))
|
|
.where(cls.manual_upload_candidate_id.is_(None))
|
|
.order_by(cls.created_at.desc())
|
|
)
|
|
rows = []
|
|
for rec, title in result.all():
|
|
job_id = rec.assigned_job_post_id or rec.job_post_id
|
|
rows.append({
|
|
"source": "form",
|
|
"email": (rec.candidate_email or "").strip().lower() or None,
|
|
"inbox_id": None,
|
|
"message_id": None,
|
|
"manual_upload_candidate_id": None,
|
|
"form_data_id": str(rec.id),
|
|
"candidate_id": None,
|
|
"job_post_id": str(job_id) if job_id else None,
|
|
"job_title": title or rec.position_applied_for or None,
|
|
"status": rec.processing_state or None,
|
|
"applied_at": (
|
|
rec.entry_date.isoformat() if rec.entry_date
|
|
else (rec.created_at.isoformat() if rec.created_at else None)
|
|
),
|
|
})
|
|
return rows
|
|
|
|
@classmethod
|
|
async def count_form_data(
|
|
cls, session: AsyncSession, *, sheet=None, search=None,
|
|
processing_state=None, is_duplicate=None,
|
|
):
|
|
statement = select(func.count()).select_from(cls)
|
|
for clause in cls._filters(
|
|
sheet=sheet, search=search,
|
|
processing_state=processing_state, is_duplicate=is_duplicate,
|
|
):
|
|
statement = statement.where(clause)
|
|
result = await session.execute(statement)
|
|
return result.scalar_one()
|
|
|
|
@classmethod
|
|
async def count_processing(cls, session: AsyncSession, *, sheet=None):
|
|
"""Tab badge counts for the Sheet Forms channel."""
|
|
statement = select(
|
|
func.count().label("all"),
|
|
func.coalesce(func.sum(case((cls.processing_state == "unread", 1), else_=0)), 0).label("unread"),
|
|
func.coalesce(func.sum(case((cls.processing_state == "imported", 1), else_=0)), 0).label("imported"),
|
|
func.coalesce(func.sum(case((cls.processing_state == "processed", 1), else_=0)), 0).label("processed"),
|
|
func.coalesce(func.sum(case((cls.processing_state == "rejected", 1), else_=0)), 0).label("rejected"),
|
|
func.coalesce(func.sum(case((cls.is_duplicate == True, 1), else_=0)), 0).label("duplicates"), # noqa: E712
|
|
).select_from(cls)
|
|
if sheet:
|
|
statement = statement.where(cls.sheet == sheet)
|
|
row = (await session.execute(statement)).one()
|
|
return {
|
|
"all": int(row.all or 0),
|
|
"unread": int(row.unread or 0),
|
|
"imported": int(row.imported or 0),
|
|
"processed": int(row.processed or 0),
|
|
"rejected": int(row.rejected or 0),
|
|
"duplicates": int(row.duplicates or 0),
|
|
}
|
|
|
|
@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}
|
|
|
|
@classmethod
|
|
async def stamp_suggested_job_posts(
|
|
cls, session: AsyncSession, records: list[dict],
|
|
) -> list[dict]:
|
|
"""Set suggested_job_post_ids from ILIKE title match on position_applied_for.
|
|
|
|
One applied-for title can match many job_posts. Blank or no match → [].
|
|
Recruiter assignment (job_post_id / assigned_job_post_id) stays unset.
|
|
"""
|
|
from job.job_post.models import JobPosts
|
|
|
|
found = await JobPosts.ids_for_titles_ilike(
|
|
session,
|
|
[r.get("position_applied_for") for r in records],
|
|
)
|
|
for record in records:
|
|
applied = (record.get("position_applied_for") or "").strip()
|
|
hits = found.get(applied) or [] if applied else []
|
|
record["suggested_job_post_ids"] = [str(uid) for uid in hits]
|
|
record["job_post_id"] = None
|
|
record["assigned_job_post_id"] = None
|
|
return records
|
|
|
|
@staticmethod
|
|
def _cell(data: dict, key: str):
|
|
"""Sheet cell → stripped str, or None if missing/blank."""
|
|
value = data.get(key)
|
|
if value is None:
|
|
return None
|
|
text = str(value).strip()
|
|
return text if text else None
|
|
|
|
@classmethod
|
|
def from_sheet_row(cls, sheet: str, row_number: int, data: dict) -> dict:
|
|
"""Build FormData kwargs from one sheet row dict (exact header keys, no aliases).
|
|
|
|
Year of Graduation: prefer the second column when present; else the first;
|
|
else None. Duplicate headers are renamed Year of Graduation_1 by normalise_headers.
|
|
"""
|
|
from g_sheet.plugins import parse_date, parse_date_time, parse_salary
|
|
|
|
first_year = cls._cell(data, "Year of Graduation")
|
|
second_year = cls._cell(data, "Year of Graduation_1")
|
|
if second_year:
|
|
entry_year = second_year
|
|
elif first_year:
|
|
entry_year = first_year
|
|
else:
|
|
entry_year = None
|
|
|
|
timestamp_raw = data.get("Timestamp")
|
|
entry_date, entry_time = parse_date_time(timestamp_raw)
|
|
|
|
current_salary = cls._cell(data, "Current Salary")
|
|
expected_salary = cls._cell(data, "Expected Salary")
|
|
|
|
return {
|
|
"sheet": sheet,
|
|
"row_number": row_number,
|
|
"raw_record": dict(data),
|
|
"entry_year": entry_year,
|
|
"entry_date": entry_date,
|
|
"entry_time": entry_time,
|
|
"name": cls._cell(data, "Full Name"),
|
|
"gender": cls._cell(data, "Gender"),
|
|
"candidate_number": cls._cell(data, "Phone number (03XX-XXXXXXX)"),
|
|
"candidate_email": cls._cell(data, "Email"),
|
|
"date_of_birth": parse_date(data.get("Date of Birth")),
|
|
"cnic": cls._cell(data, "National Identification No. (42000-XXXXXXX-X)"),
|
|
"marital_status": cls._cell(data, "Marital Status"),
|
|
"position_applied_for": cls._cell(data, "Position Applied For"),
|
|
"profile_link": cls._cell(data, "LinkedIn Profile Link"),
|
|
"residing_country": cls._cell(data, "Residing Country"),
|
|
"residing_city": cls._cell(data, "Residing City"),
|
|
"ho_availability": cls._cell(data, "Are you willing to relocate?"),
|
|
"degree": cls._cell(data, "Educational Degree"),
|
|
"university": cls._cell(data, "University"),
|
|
"university_other": cls._cell(
|
|
data,
|
|
"If your university is not listed above, please specify its name.",
|
|
),
|
|
"notice_period": cls._cell(data, "How soon can you join us?"),
|
|
"resume_link": cls._cell(data, "Drop your updated resume"),
|
|
"source_of_application": cls._cell(
|
|
data,
|
|
"Where did you hear about the position you're applying for?",
|
|
),
|
|
"cgpa": cls._cell(data, "CGPA"),
|
|
"area_of_expertise": cls._cell(data, "Area of Interest"),
|
|
"current_salary": current_salary,
|
|
"current_salary_value": parse_salary(current_salary),
|
|
"expected_salary": expected_salary,
|
|
"expected_salary_value": parse_salary(expected_salary),
|
|
"screened_by": cls._cell(data, "Recruiter"),
|
|
"hr_comments": cls._cell(data, "HR Comment"),
|
|
"director_poc_category": cls._cell(data, "Director / POC / Category"),
|
|
}
|
|
|
|
|
|
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 get_latest(cls, session: AsyncSession):
|
|
result = await session.execute(
|
|
select(cls).order_by(cls.created_at.desc()).limit(1)
|
|
)
|
|
return result.scalars().first()
|
|
|
|
@classmethod
|
|
async def delete_failed(cls, session: AsyncSession, *, commit: bool = True):
|
|
"""Drop failed import rows so a new job is not blocked by them."""
|
|
result = await session.execute(delete(cls).where(cls.status == "failed"))
|
|
if commit:
|
|
await session.commit()
|
|
return result.rowcount
|
|
|
|
@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
|