From 65c189a6dbfa52cd8fee8e2d1ac6cc5d5c93566a Mon Sep 17 00:00:00 2001 From: "ahmed.mujtaba" Date: Tue, 25 Aug 2026 20:11:02 +0500 Subject: [PATCH] . --- .dockerignore | 6 +- backend/README.md | 2 +- backend/alembic_setup.py | 10 +- backend/db_setup.py | 9 +- backend/g_sheet/app.py | 7 +- backend/g_sheet/enums.py | 359 ++++++++---------- backend/g_sheet/models.py | 118 ++++-- backend/g_sheet/plugins.py | 332 ++++++++-------- backend/g_sheet/serializers.py | 55 +-- backend/g_sheet/tasks.py | 7 +- backend/g_sheet/views.py | 23 +- backend/inbox/views.py | 17 +- backend/main.py | 13 + backend/requirements.txt | 5 + .../taskiq_management/g_sheet_broker_setup.py | 52 +++ backend/tests/test_g_sheet_plugins.py | 290 ++++++++++++++ docker-compose.dev.yml | 5 + docker-compose.yml | 18 + 18 files changed, 840 insertions(+), 488 deletions(-) create mode 100644 backend/taskiq_management/g_sheet_broker_setup.py create mode 100644 backend/tests/test_g_sheet_plugins.py diff --git a/.dockerignore b/.dockerignore index bd2cfe6..3ccd3d4 100644 --- a/.dockerignore +++ b/.dockerignore @@ -31,10 +31,8 @@ frontend/ # Candidate CVs live on the bind mount, not inside an image. backend/inbox/decoded_attachments/ -# Alembic revision scripts stay out of images (gitignored; never ship to prod). -# Schema drift is applied filelessly at API boot when DB_AUTOGENERATE=true. -backend/migrations/versions/*.py -!backend/migrations/versions/.gitkeep +# Ship revision scripts so `DB_AUTO_MIGRATE=true` can `upgrade head` in Docker. +# Fileless ORM drift (DB_AUTOGENERATE) still covers leftover model gaps. docs/ tests/ diff --git a/backend/README.md b/backend/README.md index 8811b80..c65e4b7 100644 --- a/backend/README.md +++ b/backend/README.md @@ -986,7 +986,7 @@ LLM failures are logged and skipped; the API still comes up. ```bash taskiq worker taskiq_management.broker_setup:broker \ - inbox.tasks inbox.sync_tasks taskiq_management.tasks + inbox.tasks inbox.sync_tasks taskiq_management.tasks g_sheet.tasks ``` **CV-upload worker** (isolated stream for manual uploads): diff --git a/backend/alembic_setup.py b/backend/alembic_setup.py index 4063f48..4d9ec35 100644 --- a/backend/alembic_setup.py +++ b/backend/alembic_setup.py @@ -211,7 +211,10 @@ def context_options() -> dict[str, Any]: async def _run(fn: Callable[[Connection], Any]) -> Any: """Run a synchronous Alembic call on the async engine's connection.""" + schema = get_settings().db_default_schema or "public" async with get_engine().connect() as conn: + # Unqualified FK targets (REFERENCES users) must resolve in `app`. + await conn.execute(text(f'SET search_path TO "{schema}", public')) result = await conn.run_sync(fn) await conn.commit() return result @@ -418,7 +421,12 @@ async def migrate(*, autogen: bool | None = None, message: str = "auto") -> None else: await upgrade() if should_autogen: - await apply_model_drift() + try: + await apply_model_drift() + except Exception as exc: + # Drift can still trip on unrelated tables; file revisions already + # ran above. Log and continue so the API can finish booting. + logger.exception("ORM drift apply failed; continuing boot: %s", exc) await run_manual_sql() logger.info("database at revision %s", await current()) diff --git a/backend/db_setup.py b/backend/db_setup.py index 87b38f0..811cff2 100644 --- a/backend/db_setup.py +++ b/backend/db_setup.py @@ -174,8 +174,15 @@ def _connect_args(settings: Settings) -> dict: """UTC session + SSL for RDS. `require` encrypts without verifying the CA.""" import ssl as ssl_mod + # search_path includes the app schema so unqualified FKs (users.id) resolve + # during fileless ORM drift and normal queries — default is "$user", public. + schema = settings.db_default_schema or "public" args: dict = { - "server_settings": {"timezone": "UTC", "application_name": settings.app_name} + "server_settings": { + "timezone": "UTC", + "application_name": settings.app_name, + "search_path": f"{schema}, public", + } } mode = (settings.db_sslmode or "").strip().lower() if mode and mode not in ("disable", "allow", "prefer"): diff --git a/backend/g_sheet/app.py b/backend/g_sheet/app.py index 07f8769..f446913 100644 --- a/backend/g_sheet/app.py +++ b/backend/g_sheet/app.py @@ -100,13 +100,14 @@ async def fetch_sheet( @router.post("/sheet/import") async def import_all_sheets( + tab: str | None = Query(None), current_user: dict = Depends(require_permission(PermissionTag.SETTINGS_EDIT)), session: AsyncSession = Depends(get_session), ): - """Enqueue a full-spreadsheet import. Poll GET /sheet/import/fetch for status.""" + """No tab -> every tab. With a tab -> that sheet only. Poll GET /sheet/import/fetch.""" try: service=Sheet(session=session) - data=await service.start_import(current_user=current_user,tab=None) + data=await service.start_import(current_user=current_user,tab=tab) return JSONResponse(content={"data":data,"total":1,"status_code":200}) except HTTPException: raise @@ -179,7 +180,7 @@ async def fetch_form_data( raise except Exception as e: raise HTTPException(status_code=500,detail=str(e)) - + @router.get("/sheet/form-data/{record_id}") async def fetch_form_data_by_id( diff --git a/backend/g_sheet/enums.py b/backend/g_sheet/enums.py index 0232c91..0db456a 100644 --- a/backend/g_sheet/enums.py +++ b/backend/g_sheet/enums.py @@ -1,207 +1,193 @@ -"""Sheet header aliases, FormData keys, and date/round format mappings. +"""Sheet header aliases, FormData keys, and date format mappings. (str, Enum) like inbox/enums.py: members compare to and serialize as plain strings. -Non-string mappings (month pairs, ordinal slot+pattern) use plain Enum. +Non-string mappings (month pairs) use plain Enum. """ from enum import Enum +class AliasEnum(str, Enum): + """Member-less base so alias enums share one `has` without 25 copies.""" + + @classmethod + def has(cls, value) -> bool: + return value in cls._value2member_map_ + + class FormDataField(str, Enum): - """Canonical FormData column keys (plus title, which stays in JSONB only).""" + """Canonical FormData column keys for the recruitment screening sheet.""" + SERIAL_NO = "serial_no" + ENTRY_YEAR = "entry_year" + ENTRY_MONTH = "entry_month" + ENTRY_DATE = "entry_date" + ENTRY_TIME = "entry_time" + SCREENED_BY = "screened_by" NAME = "name" - DEGREE = "degree" - EXPERIENCE = "experience" + HR_COMMENTS = "hr_comments" + CANDIDATE_NUMBER = "candidate_number" + CANDIDATE_EMAIL = "candidate_email" + PROFILE_LINK = "profile_link" + AREA_OF_EXPERTISE = "area_of_expertise" + REQUISITION_NUMBER = "requisition_number" + POSITION_SUITABLE_FOR = "position_suitable_for" + SOURCE_OF_APPLICATION = "source_of_application" AGE = "age" - FAMILY_DETAILS = "family_details" - TITLE = "title" - - -class NameAlias(str, Enum): - NAME = "name" - NAMES = "names" - CANDIDATE_NAME = "candidate name" - CANDIDATE = "candidate" - - @classmethod - def has(cls, value) -> bool: - return value in cls._value2member_map_ - - -class DegreeAlias(str, Enum): - EDUCATION = "education" + MARITAL_STATUS = "marital_status" DEGREE = "degree" - QUALIFICATION = "qualification" - - @classmethod - def has(cls, value) -> bool: - return value in cls._value2member_map_ - - -class ExperienceAlias(str, Enum): + UNIVERSITY = "university" EXPERIENCE = "experience" - EXP = "exp" - YEARS_OF_EXPERIENCE = "years of experience" - TOTAL_EXPERIENCE = "total experience" - - @classmethod - def has(cls, value) -> bool: - return value in cls._value2member_map_ + EXPERIENCE_DETAILS = "experience_details" + AREA_OF_RESIDENCE = "area_of_residence" + COMMUNICATION_SKILLS = "communication_skills" + PREFERRED_TIMINGS = "preferred_timings" + HO_AVAILABILITY = "ho_availability" + CURRENT_COMPANY = "current_company" + REASON_FOR_LEAVING = "reason_for_leaving" + NOTICE_PERIOD = "notice_period" + CURRENT_SALARY = "current_salary" + EXPECTED_SALARY = "expected_salary" + PROS = "pros" + CONS = "cons" -class AgeAlias(str, Enum): - AGE = "age" - - @classmethod - def has(cls, value) -> bool: - return value in cls._value2member_map_ - - -class FamilyDetailsAlias(str, Enum): - FAMILY_DETAILS = "family details" - MARITAL_STATUS = "marital status" - MARITAL = "marital" - - @classmethod - def has(cls, value) -> bool: - return value in cls._value2member_map_ - - -class TitleAlias(str, Enum): - """No FormData column — recognised so headers are not treated as unknown noise.""" - - TITLE = "title" - DESIGNATION = "designation" - ROLE = "role" - POSITION = "position" - TEAM = "team" - JOB_TITLE = "job title" - AREA_OF_EXPERTISE = "area of expertise" - DEPARTMENT = "department" - - @classmethod - def has(cls, value) -> bool: - return value in cls._value2member_map_ - - -# FormDataField → alias Enum. Order is match priority for overlapping startswith hits. -FIELD_ALIAS_ENUMS = { - FormDataField.NAME: NameAlias, - FormDataField.DEGREE: DegreeAlias, - FormDataField.EXPERIENCE: ExperienceAlias, - FormDataField.AGE: AgeAlias, - FormDataField.FAMILY_DETAILS: FamilyDetailsAlias, - FormDataField.TITLE: TitleAlias, +# canonical (lowercased, whitespace-collapsed, punctuation-stripped) header -> field. +# The sheet's own spelling is listed first; the rest are tolerated synonyms. +HEADER_ALIASES: dict[FormDataField, tuple[str, ...]] = { + FormDataField.SERIAL_NO: ("um", "sr", "sr no", "s no", "serial", "serial no"), + FormDataField.ENTRY_YEAR: ("year", "year of graduation"), + FormDataField.ENTRY_MONTH: ("month",), + FormDataField.ENTRY_DATE: ("date", "entry date", "date of entry", "timestamp", "time stamp"), + FormDataField.ENTRY_TIME: ("time of entry", "entry time", "time"), + FormDataField.SCREENED_BY: ( + "screened by", "screened", "interviewed by", "conducted by", "recruiter", + ), + FormDataField.NAME: ( + "candidate name", "full name", "name", "names", "candidate", + ), + FormDataField.HR_COMMENTS: ("hr comments", "hr comment", "comments", "remarks"), + FormDataField.CANDIDATE_NUMBER: ( + "candidate number", "contact number", "phone number", "phone", "mobile", "contact", + ), + FormDataField.CANDIDATE_EMAIL: ("candidate email", "email", "email address"), + FormDataField.PROFILE_LINK: ( + "profile link", "linkedin profile link", "cv link", "resume link", + "drop your updated resume", "profile", + ), + FormDataField.AREA_OF_EXPERTISE: ( + "area of expertise", "area of interest", "expertise", + ), + FormDataField.REQUISITION_NUMBER: ("requisition number", "requisition", "req no"), + FormDataField.POSITION_SUITABLE_FOR: ( + "position suitable for", "position applied for", "position", + "designation", "job title", "role", "title", + ), + FormDataField.SOURCE_OF_APPLICATION: ( + "source of application", "source", "application source", + "where did you hear about the position you're applying for", + ), + FormDataField.AGE: ("age",), + FormDataField.MARITAL_STATUS: ("marital status", "marital", "family details"), + FormDataField.DEGREE: ( + "education", "educational degree", "degree", "qualification", + ), + FormDataField.UNIVERSITY: ("university of graduation", "university", "institute", "college"), + FormDataField.EXPERIENCE: ("experience", "total experience", "years of experience", "exp"), + FormDataField.EXPERIENCE_DETAILS: ("experience details", "experience detail"), + FormDataField.AREA_OF_RESIDENCE: ( + "area of residence", "residing city", "residing country", + "residence", "location", "address", + ), + FormDataField.COMMUNICATION_SKILLS: ("communication skills", "communication"), + FormDataField.PREFERRED_TIMINGS: ("preferred timings", "preferred timing", "shift"), + FormDataField.HO_AVAILABILITY: ( + "availability to work in the h.o", "availability to work in the ho", + "ho availability", "availability", "are you willing to relocate", + ), + FormDataField.CURRENT_COMPANY: ("current company", "current employer", "company", "employer"), + FormDataField.REASON_FOR_LEAVING: ("reason for leaving", "reason of leaving", "reason"), + FormDataField.NOTICE_PERIOD: ( + "how soon can you join us", "how soon can you join", + "notice period", "joining", "availability to join", + ), + FormDataField.CURRENT_SALARY: ("current salary", "present salary", "salary"), + FormDataField.EXPECTED_SALARY: ("expected salary", "salary expectation", "expected"), + FormDataField.PROS: ("pros", "strengths"), + FormDataField.CONS: ("cons", "weaknesses"), } -class RoundRole(str, Enum): - """Interview-round column roles resolved left-to-right into four slots.""" - - DATE = "date" - BY = "by" - STATUS = "status" - NOTES = "notes" - RESULT = "result" +def _build_alias_to_field() -> dict[str, FormDataField]: + inverted: dict[str, FormDataField] = {} + for field, aliases in HEADER_ALIASES.items(): + for alias in aliases: + if alias in inverted: + raise ValueError( + f"duplicate header alias {alias!r}: " + f"{inverted[alias].value} and {field.value}" + ) + inverted[alias] = field + return inverted -class ConductedByAlias(str, Enum): - """Header spellings that map to RoundRole.BY.""" - - CONDUCTED_BY = "conducted by" - INTERVIEWED_BY = "interviewed by" - INTERVIEW_BY = "interview by" - CONDUCTED = "conducted" - BY = "by" - - @classmethod - def has(cls, value) -> bool: - return value in cls._value2member_map_ - - @classmethod - def contained_in(cls, text: str) -> bool: - return any(member.value in text for member in cls if " " in member.value) +ALIAS_TO_FIELD: dict[str, FormDataField] = _build_alias_to_field() -class NotesToken(str, Enum): - """Substrings that classify a header as RoundRole.NOTES.""" +class FormDataColumn(str, Enum): + """FormData API / ORM field names in serialize order. - NOTE = "note" - REMARK = "remark" - COMMENT = "comment" + Broader than FormDataField: includes id, sheet meta, derived parsers + (age_raw, *_salary_value), raw_record, and timestamps. + """ - @classmethod - def contained_in(cls, text: str) -> bool: - return any(member.value in text for member in cls) + ID = "id" + SHEET = "sheet" + JOB_POST_ID = "job_post_id" + ROW_NUMBER = "row_number" + SERIAL_NO = "serial_no" + ENTRY_YEAR = "entry_year" + ENTRY_MONTH = "entry_month" + ENTRY_DATE = "entry_date" + ENTRY_TIME = "entry_time" + SCREENED_BY = "screened_by" + NAME = "name" + HR_COMMENTS = "hr_comments" + CANDIDATE_NUMBER = "candidate_number" + CANDIDATE_EMAIL = "candidate_email" + PROFILE_LINK = "profile_link" + AREA_OF_EXPERTISE = "area_of_expertise" + REQUISITION_NUMBER = "requisition_number" + POSITION_SUITABLE_FOR = "position_suitable_for" + SOURCE_OF_APPLICATION = "source_of_application" + AGE = "age" + AGE_RAW = "age_raw" + MARITAL_STATUS = "marital_status" + DEGREE = "degree" + UNIVERSITY = "university" + EXPERIENCE = "experience" + EXPERIENCE_DETAILS = "experience_details" + AREA_OF_RESIDENCE = "area_of_residence" + COMMUNICATION_SKILLS = "communication_skills" + PREFERRED_TIMINGS = "preferred_timings" + HO_AVAILABILITY = "ho_availability" + CURRENT_COMPANY = "current_company" + REASON_FOR_LEAVING = "reason_for_leaving" + NOTICE_PERIOD = "notice_period" + CURRENT_SALARY = "current_salary" + CURRENT_SALARY_VALUE = "current_salary_value" + EXPECTED_SALARY = "expected_salary" + EXPECTED_SALARY_VALUE = "expected_salary_value" + PROS = "pros" + CONS = "cons" + RAW_RECORD = "raw_record" + IMPORTED_AT = "imported_at" + CREATED_AT = "created_at" + UPDATED_AT = "updated_at" -# -- Round → FormData column names (slot 0..3 = definition order) ------------ - -class RoundDateColumn(str, Enum): - R1 = "interview_date" - R2 = "second_interview_date" - R3 = "third_interview_date" - R4 = "fourth_interview_date" - - @classmethod - def ordered(cls) -> tuple[str, ...]: - return tuple(member.value for member in cls) - - -class RoundByColumn(str, Enum): - R1 = "interview_by" - R2 = "second_interview_by" - R3 = "third_interview_by" - R4 = "fourth_interview_by" - - @classmethod - def ordered(cls) -> tuple[str, ...]: - return tuple(member.value for member in cls) - - -class RoundTimeColumn(str, Enum): - R1 = "interview_time" - R2 = "second_interview_time" - R3 = "third_interview_time" - R4 = "fourth_interview_time" - - @classmethod - def ordered(cls) -> tuple[str, ...]: - return tuple(member.value for member in cls) - - -class RoundStatusColumn(str, Enum): - R1 = "interview_status" - R2 = "second_interview_status" - R3 = "third_interview_status" - R4 = "fourth_interview_status" - - @classmethod - def ordered(cls) -> tuple[str, ...]: - return tuple(member.value for member in cls) - - -class RoundNotesColumn(str, Enum): - R1 = "interview_notes" - R2 = "second_interview_notes" - R3 = "third_interview_notes" - R4 = "fourth_interview_notes" - - @classmethod - def ordered(cls) -> tuple[str, ...]: - return tuple(member.value for member in cls) - - -class RoundResultColumn(str, Enum): - R1 = "interview_result" - R2 = "second_interview_result" - R3 = "third_interview_result" - R4 = "fourth_interview_result" - - @classmethod - def ordered(cls) -> tuple[str, ...]: - return tuple(member.value for member in cls) +# Ordered values for serialize_form_data / model_fields assertions. +FORM_DATA_FIELDS: tuple[str, ...] = tuple(member.value for member in FormDataColumn) # -- Date parsing ------------------------------------------------------------ @@ -260,20 +246,3 @@ class MonthNormalisation(Enum): @property def short(self) -> str: return self.value[1] - - -class RoundOrdinal(Enum): - """Interview-round ordinal in a header → slot index 0..3. value is (slot, regex).""" - - FIRST = (0, r"(?:1st|first|01st)") - SECOND = (1, r"(?:2nd|second|02nd)") - THIRD = (2, r"(?:3rd|third|03rd)") - FOURTH = (3, r"(?:4th|fourth|04th)") - - @property - def slot(self) -> int: - return self.value[0] - - @property - def pattern(self) -> str: - return self.value[1] diff --git a/backend/g_sheet/models.py b/backend/g_sheet/models.py index 502e0c8..e63703c 100644 --- a/backend/g_sheet/models.py +++ b/backend/g_sheet/models.py @@ -5,7 +5,7 @@ from __future__ import annotations import uuid from datetime import datetime, timezone -from sqlalchemy import Column, DateTime, Index, delete, func, or_ +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 @@ -15,6 +15,9 @@ 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.""" @@ -23,45 +26,49 @@ class FormData(SQLModel, table=True): Index("ix_form_data_sheet_row_number", "sheet", "row_number", unique=True), ) - id: int | None = Field(default=None, primary_key=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) - degree: str | None = Field(default=None) - experience: 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) + 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) - family_details: str | None = Field(default=None) - - interview_date: datetime | None = Field(default=None, sa_type=DateTime(timezone=True)) - interview_by: str | None = Field(default=None) - interview_time: str | None = Field(default=None) - interview_status: str | None = Field(default=None) - interview_notes: str | None = Field(default=None) - interview_result: str | None = Field(default=None) - - second_interview_date: datetime | None = Field(default=None, sa_type=DateTime(timezone=True)) - second_interview_by: str | None = Field(default=None) - second_interview_time: str | None = Field(default=None) - second_interview_status: str | None = Field(default=None) - second_interview_notes: str | None = Field(default=None) - second_interview_result: str | None = Field(default=None) - - third_interview_date: datetime | None = Field(default=None, sa_type=DateTime(timezone=True)) - third_interview_by: str | None = Field(default=None) - third_interview_time: str | None = Field(default=None) - third_interview_status: str | None = Field(default=None) - third_interview_notes: str | None = Field(default=None) - third_interview_result: str | None = Field(default=None) - - fourth_interview_date: datetime | None = Field(default=None, sa_type=DateTime(timezone=True)) - fourth_interview_by: str | None = Field(default=None) - fourth_interview_time: str | None = Field(default=None) - fourth_interview_status: str | None = Field(default=None) - fourth_interview_notes: str | None = Field(default=None) - fourth_interview_result: str | None = Field(default=None) + marital_status: str | None = Field(default=None) + degree: str | None = Field(default=None) + university: 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) + 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) + pros: str | None = Field(default=None) + cons: str | None = Field(default=None) raw_record: dict | None = Field(default=None, sa_column=Column(JSONB)) - row_number: int | None = Field(default=None) 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)) @@ -72,19 +79,30 @@ class FormData(SQLModel, table=True): 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.interview_by.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), )) return filters @classmethod async def get_form_data_by_id(cls, session: AsyncSession, record_id): try: - rid = int(record_id) + rid = uuid.UUID(str(record_id)) except (TypeError, ValueError): return None result = await session.execute(select(cls).where(cls.id == rid)) @@ -129,12 +147,28 @@ class FormData(SQLModel, table=True): return deleted @classmethod - async def insert_form_data_bulk(cls, session: AsyncSession, records: list[dict], *, commit: bool = True): - rows = [cls(**fields) for fields in records] - session.add_all(rows) + 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 len(rows) + return total @classmethod async def replace_sheet(cls, session: AsyncSession, sheet: str, records: list[dict]): @@ -153,7 +187,9 @@ class SheetImportRun(SQLModel, table=True): 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) - created_by: uuid.UUID | None = Field(default=None, foreign_key="users.id") + # 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) diff --git a/backend/g_sheet/plugins.py b/backend/g_sheet/plugins.py index 34001b2..9d6645e 100644 --- a/backend/g_sheet/plugins.py +++ b/backend/g_sheet/plugins.py @@ -24,21 +24,11 @@ from googleapiclient.discovery import build from googleapiclient.errors import HttpError from g_sheet.enums import ( - ConductedByAlias, + ALIAS_TO_FIELD, DateFormat, DateTimeSeparator, - FIELD_ALIAS_ENUMS, FormDataField, MonthNormalisation, - NotesToken, - RoundByColumn, - RoundDateColumn, - RoundNotesColumn, - RoundOrdinal, - RoundResultColumn, - RoundRole, - RoundStatusColumn, - RoundTimeColumn, ) load_dotenv() @@ -223,18 +213,27 @@ def rows_to_records(rows): Sheets truncates trailing empties, so short rows are padded to header width. Fully blank rows are dropped rather than emitted as all-empty records. """ + return [record for _,record in rows_to_indexed_records(rows)] + + +def rows_to_indexed_records(rows): + """Sheet rows -> (1-based sheet row number, record) pairs. + + Blank interior rows are skipped but do not shift later row numbers — the index + is the true sheet row (header is row 1), which is half of the unique key. + """ if not rows: return [] headers=normalise_headers(rows[0]) - records=[] - for row in rows[1:]: + indexed=[] + for offset,row in enumerate(rows[1:]): values=[str(cell) if cell is not None else "" for cell in row] if not any(value.strip() for value in values): continue if len(values)\d+(?:[.,]\d+)?)\s*(?Pk|lac|lakh|lacs|lakhs|crore|crores)?\b", + re.I, +) +_CURRENCY_STRIP_RE=re.compile(r"(?:rs\.?|pkr|inr|usd|\$|€|£)",re.I) + +# Every typed column key the mapper must emit (uniform dicts for bulk insert). +_FORM_DATA_COLUMN_KEYS=tuple(field.value for field in FormDataField)+( + "age_raw","current_salary_value","expected_salary_value", +) def canonical_header(h): - """Lower, collapse whitespace (incl. embedded newlines), strip _N and (tails).""" + """Lower, collapse whitespace (incl. embedded newlines), strip _N, (tails), trailing punct.""" text=str(h or "").replace("\n"," ").replace("\r"," ") text=re.sub(r"\s+"," ",text).strip().lower() text=re.sub(r"_\d+$","",text) text=re.sub(r"\s*\([^)]*\)\s*$","",text).strip() + text=text.rstrip("?:.,").strip() return text def match_field(h): - """Map a sheet header to a FormDataField, or None. - - Exact alias first, then startswith. No fuzzy matching — dirty headers mislabel - more often than they rescue, and a miss is non-fatal (value stays in JSONB). - """ + """Map a sheet header to a FormDataField via exact alias lookup, or None.""" canon=canonical_header(h) if not canon: return None - for field,alias_enum in FIELD_ALIAS_ENUMS.items(): - if alias_enum.has(canon): - return field - for field,alias_enum in FIELD_ALIAS_ENUMS.items(): - for alias in alias_enum: - if canon.startswith(alias.value): - return field - return None + return ALIAS_TO_FIELD.get(canon) def resolve_name(record,headers): - """Candidate name: alias match, else column A (headers[0]) — always the name.""" + """Candidate name: alias match, else first non-meta column (not Timestamp/date).""" for header in headers: if match_field(header)==FormDataField.NAME: value=record.get(header) if value is not None and str(value).strip(): return str(value).strip() - if headers: - value=record.get(headers[0]) + # Skip entry/meta columns so Google Form "Timestamp" is never treated as a name. + _skip={ + FormDataField.ENTRY_DATE,FormDataField.ENTRY_TIME, + FormDataField.ENTRY_YEAR,FormDataField.ENTRY_MONTH,FormDataField.SERIAL_NO, + } + for header in headers: + if match_field(header) in _skip: + continue + value=record.get(header) if value is not None and str(value).strip(): return str(value).strip() return None -def _classify_round_role(canon): - """RoundRole for a canonical header, or None for unrecognised headers.""" - if not canon: - return None - if ConductedByAlias.contained_in(canon) or ConductedByAlias.has(canon): - return RoundRole.BY - if canon.startswith(ConductedByAlias.CONDUCTED.value): - return RoundRole.BY - if RoundRole.RESULT.value in canon: - return RoundRole.RESULT - if RoundRole.STATUS.value in canon: - return RoundRole.STATUS - if NotesToken.contained_in(canon): - return RoundRole.NOTES - if RoundRole.DATE.value in canon: - return RoundRole.DATE - return None - - -def _extract_ordinal(canon): - for slot,pattern in _ORDINAL_PATTERNS: - if pattern.search(canon): - return slot - return None - - -def resolve_round_columns(headers): - """Positional interview-round map: scan left→right into four slots. - - Ordinal in the header (`2nd`, `second`) pins the slot; otherwise the first free - slot for that role is taken, never moving backwards. A fifth Results_4 stays - unmapped (JSONB). Literal-date headers like `19-Feb-2026` classify as nothing. - """ - slots=[{role:None for role in RoundRole} for _ in range(4)] - cursor={role:0 for role in RoundRole} - - for header in headers: - canon=canonical_header(header) - role=_classify_round_role(canon) - if role is None: - continue - ordinal=_extract_ordinal(canon) - if ordinal is not None: - if slots[ordinal][role] is None: - slots[ordinal][role]=header - continue - start=cursor[role] - chosen=None - for index in range(start,4): - if slots[index][role] is None: - chosen=index - break - if chosen is None: - continue - slots[chosen][role]=header - cursor[role]=chosen+1 - return slots - - def _normalise_month_spellings(text): """strptime %b rejects `Sept`; expand common sheet spellings first.""" lowered=text.lower() @@ -404,7 +340,7 @@ def parse_date(value): def parse_date_time(value): - """(datetime|None, time_string|None) — fills *_time for the cells that carry one.""" + """(datetime|None, time_string|None) — fills entry_time when the cell carries one.""" parsed=parse_date(value) if value is None: return parsed,None @@ -430,6 +366,53 @@ def parse_age(value): return None,raw +def parse_score(value): + """First digit run kept only when 0 <= n <= 10 (communication skills scale).""" + if value is None: + return None + text=str(value).strip() + if not text: + return None + match=_SCORE_RE.search(text) + if not match: + return None + number=int(match.group()) + if 0<=number<=10: + return number + return None + + +def parse_salary(value): + """Numeric salary in whole currency units, or None for non-numeric cells. + + Understands k/K, lac/lakh, crore; on a range takes the first number. + The raw cell text still goes to *_salary — a None here loses nothing. + """ + if value is None: + return None + text=str(value).strip() + if not text: + return None + cleaned=_CURRENCY_STRIP_RE.sub(" ",text) + cleaned=cleaned.replace(",","") + match=_SALARY_UNIT_RE.search(cleaned) + if not match: + return None + raw_num=match.group("num").replace(",","") + try: + amount=float(raw_num) + except ValueError: + return None + unit=(match.group("unit") or "").lower() + if unit=="k": + amount*=1000 + elif unit in ("lac","lakh","lacs","lakhs"): + amount*=100000 + elif unit in ("crore","crores"): + amount*=10000000 + return int(amount) + + def _blank_to_none(value): if value is None: return None @@ -437,100 +420,95 @@ def _blank_to_none(value): return text if text else None -def map_record_to_form_data(sheet,record,headers,row_number): - """Pure row mapper → kwargs dict for FormData(**...).""" - rounds=resolve_round_columns(headers) - mapped={ - "sheet":sheet, - "row_number":row_number, - "raw_record":dict(record), - "name":_blank_to_none(resolve_name(record,headers)), - "degree":None, - "experience":None, - "age":None, - "age_raw":None, - "family_details":None, - } - for field in BY_FIELDS+TIME_FIELDS+STATUS_FIELDS+NOTES_FIELDS+RESULT_FIELDS: - mapped[field]=None - for field in DATE_FIELDS: - mapped[field]=None - - for header,value in record.items(): +def _header_field_map(headers): + """header -> FormDataField, first header that claims each field wins.""" + claimed={} + header_to_field={} + for header in headers: field=match_field(header) - if field==FormDataField.DEGREE: - mapped["degree"]=_blank_to_none(value) - elif field==FormDataField.EXPERIENCE: - mapped["experience"]=_blank_to_none(value) - elif field==FormDataField.AGE: + if field is None or field in claimed: + continue + claimed[field]=header + header_to_field[header]=field + return header_to_field + + +def map_record_to_form_data(sheet,record,headers,row_number): + """Pure row mapper → kwargs dict for FormData (uniform keys for bulk insert).""" + mapped={key:None for key in _FORM_DATA_COLUMN_KEYS} + mapped["sheet"]=sheet + mapped["row_number"]=row_number + mapped["raw_record"]=dict(record) + mapped["name"]=_blank_to_none(resolve_name(record,headers)) + + for header,field in _header_field_map(headers).items(): + value=record.get(header) + key=field.value + if field==FormDataField.AGE: age,age_raw=parse_age(value) mapped["age"]=age mapped["age_raw"]=age_raw - elif field==FormDataField.FAMILY_DETAILS: - mapped["family_details"]=_blank_to_none(value) - - for index,slot in enumerate(rounds): - if slot.get(RoundRole.DATE): - dt,tm=parse_date_time(record.get(slot[RoundRole.DATE])) - mapped[DATE_FIELDS[index]]=dt - mapped[TIME_FIELDS[index]]=tm - if slot.get(RoundRole.BY): - mapped[BY_FIELDS[index]]=_blank_to_none(record.get(slot[RoundRole.BY])) - if slot.get(RoundRole.STATUS): - mapped[STATUS_FIELDS[index]]=_blank_to_none(record.get(slot[RoundRole.STATUS])) - if slot.get(RoundRole.NOTES): - mapped[NOTES_FIELDS[index]]=_blank_to_none(record.get(slot[RoundRole.NOTES])) - if slot.get(RoundRole.RESULT): - mapped[RESULT_FIELDS[index]]=_blank_to_none(record.get(slot[RoundRole.RESULT])) + elif field==FormDataField.ENTRY_DATE: + dt,tm=parse_date_time(value) + mapped["entry_date"]=dt + if tm and not mapped.get("entry_time"): + mapped["entry_time"]=tm + elif field==FormDataField.COMMUNICATION_SKILLS: + mapped["communication_skills"]=parse_score(value) + elif field==FormDataField.CURRENT_SALARY: + mapped["current_salary"]=_blank_to_none(value) + mapped["current_salary_value"]=parse_salary(value) + elif field==FormDataField.EXPECTED_SALARY: + mapped["expected_salary"]=_blank_to_none(value) + mapped["expected_salary_value"]=parse_salary(value) + elif field==FormDataField.NAME: + # resolve_name already set this; keep its column-A fallback behaviour. + continue + else: + mapped[key]=_blank_to_none(value) return mapped def collect_unmapped_headers(headers): - """Headers that are neither a typed alias nor claimed by a round slot. - - `title` aliases are included — they have no FormData column and live in JSONB. - """ - rounds=resolve_round_columns(headers) - claimed=set() - for slot in rounds: - for role in RoundRole: - if slot.get(role): - claimed.add(slot[role]) - unmapped=[] - for header in headers: - if header in claimed: - continue - field=match_field(header) - if field is None or field==FormDataField.TITLE: - unmapped.append(header) - return unmapped + """Headers that do not exact-match any alias.""" + return [header for header in headers if match_field(header) is None] def import_row_stats(mapped_rows,headers): """Aggregate parse diagnostics for an import report.""" + unmapped=collect_unmapped_headers(headers) dates_parsed=0 dates_unparsed=0 ages_parsed=0 + salaries_parsed=0 + # Find which raw header feeds entry_date (if any) once, not per row. + entry_date_header=None + for header in headers: + if match_field(header)==FormDataField.ENTRY_DATE: + entry_date_header=header + break for row in mapped_rows: - raw=row.get("raw_record") or {} - rounds=resolve_round_columns(headers) - for index,slot in enumerate(rounds): - header=slot.get(RoundRole.DATE) - if not header: - continue - cell=raw.get(header) - if cell is None or not str(cell).strip(): - continue - if row.get(DATE_FIELDS[index]) is not None: - dates_parsed+=1 - elif _DIGIT_RE.search(str(cell)): - dates_unparsed+=1 + if entry_date_header is not None: + raw=row.get("raw_record") or {} + cell=raw.get(entry_date_header) + if cell is not None and str(cell).strip(): + if row.get("entry_date") is not None: + dates_parsed+=1 + elif _DIGIT_RE.search(str(cell)): + dates_unparsed+=1 if row.get("age") is not None: ages_parsed+=1 + if ( + row.get("current_salary_value") is not None + or row.get("expected_salary_value") is not None + ): + salaries_parsed+=1 return { "dates_parsed":dates_parsed, "dates_unparsed":dates_unparsed, "ages_parsed":ages_parsed, - "unmapped_headers":collect_unmapped_headers(headers), + "salaries_parsed":salaries_parsed, + "unmapped_headers":unmapped, } + diff --git a/backend/g_sheet/serializers.py b/backend/g_sheet/serializers.py index c74fd82..affd4e2 100644 --- a/backend/g_sheet/serializers.py +++ b/backend/g_sheet/serializers.py @@ -2,6 +2,11 @@ from __future__ import annotations +import uuid +from datetime import datetime + +from g_sheet.enums import FORM_DATA_FIELDS + def serialize_metadata(payload: dict) -> dict: """spreadsheets.get response -> the spreadsheet header the UI renders.""" @@ -99,45 +104,16 @@ def _iso(value): def serialize_form_data(row) -> dict: """FormData ORM row → API dict, including raw_record.""" - return { - "id": row.id, - "sheet": row.sheet, - "name": row.name, - "degree": row.degree, - "experience": row.experience, - "age": row.age, - "age_raw": row.age_raw, - "family_details": row.family_details, - "interview_date": _iso(row.interview_date), - "interview_by": row.interview_by, - "interview_time": row.interview_time, - "interview_status": row.interview_status, - "interview_notes": row.interview_notes, - "interview_result": row.interview_result, - "second_interview_date": _iso(row.second_interview_date), - "second_interview_by": row.second_interview_by, - "second_interview_time": row.second_interview_time, - "second_interview_status": row.second_interview_status, - "second_interview_notes": row.second_interview_notes, - "second_interview_result": row.second_interview_result, - "third_interview_date": _iso(row.third_interview_date), - "third_interview_by": row.third_interview_by, - "third_interview_time": row.third_interview_time, - "third_interview_status": row.third_interview_status, - "third_interview_notes": row.third_interview_notes, - "third_interview_result": row.third_interview_result, - "fourth_interview_date": _iso(row.fourth_interview_date), - "fourth_interview_by": row.fourth_interview_by, - "fourth_interview_time": row.fourth_interview_time, - "fourth_interview_status": row.fourth_interview_status, - "fourth_interview_notes": row.fourth_interview_notes, - "fourth_interview_result": row.fourth_interview_result, - "raw_record": row.raw_record, - "row_number": row.row_number, - "imported_at": _iso(row.imported_at), - "created_at": _iso(row.created_at), - "updated_at": _iso(row.updated_at), - } + out = {} + for key in FORM_DATA_FIELDS: + value = getattr(row, key) + if isinstance(value, datetime): + out[key] = _iso(value) + elif isinstance(value, uuid.UUID): + out[key] = str(value) + else: + out[key] = value + return out def serialize_import(report: dict) -> dict: @@ -150,6 +126,7 @@ def serialize_import(report: dict) -> dict: "dates_parsed": report.get("dates_parsed", 0), "dates_unparsed": report.get("dates_unparsed", 0), "ages_parsed": report.get("ages_parsed", 0), + "salaries_parsed": report.get("salaries_parsed", 0), "unmapped_headers": report.get("unmapped_headers") or [], "error": report.get("error"), } diff --git a/backend/g_sheet/tasks.py b/backend/g_sheet/tasks.py index 488b2d5..619d4d2 100644 --- a/backend/g_sheet/tasks.py +++ b/backend/g_sheet/tasks.py @@ -1,4 +1,4 @@ -"""Google Sheet → FormData import Taskiq tasks (shared inbox worker stream).""" +"""Google Sheet → FormData import Taskiq tasks (dedicated sheet_import stream).""" from __future__ import annotations @@ -12,7 +12,8 @@ from dotenv import load_dotenv from db_setup import session_scope from g_sheet.models import SheetImportRun from g_sheet.views import Sheet -from taskiq_management.broker_setup import MAX_RETRIES,RETRY_DELAY,broker +from taskiq_management.broker_setup import MAX_RETRIES,RETRY_DELAY +from taskiq_management.g_sheet_broker_setup import sheet_broker from taskiq_management.middleware import PermanentTaskError load_dotenv() @@ -33,7 +34,7 @@ async def _fail(run_id:str,error:str) -> dict: return {"status":"failed","error":error} -@broker.task( +@sheet_broker.task( task_name="g_sheet.import_sheets", retry_on_error=True, max_retries=MAX_RETRIES, diff --git a/backend/g_sheet/views.py b/backend/g_sheet/views.py index 96bf700..ef14797 100644 --- a/backend/g_sheet/views.py +++ b/backend/g_sheet/views.py @@ -24,7 +24,9 @@ from g_sheet.plugins import ( import_row_stats, load_credentials, map_record_to_form_data, + normalise_headers, quote_tab, + rows_to_indexed_records, rows_to_records, stringify_rows, ) @@ -202,17 +204,21 @@ class Sheet: 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_records(tab) - records=data["records"] - headers=data["headers"] - mapped=[] - for index,record in enumerate(records): - mapped.append(map_record_to_form_data(tab,record,headers,index+2)) + data=await self.read_range(tab) + rows=data["rows"] + if not rows: + return serialize_import({"tab":tab,"rows_read":0,"inserted":0,"deleted":0}) + headers=normalise_headers(rows[0]) + indexed=rows_to_indexed_records(rows) + mapped=[ + map_record_to_form_data(tab,record,headers,row_number) + for row_number,record in indexed + ] result=await FormData.replace_sheet(session,tab,mapped) stats=import_row_stats(mapped,headers) return serialize_import({ "tab":tab, - "rows_read":len(records), + "rows_read":len(indexed), "inserted":result["inserted"], "deleted":result["deleted"], **stats, @@ -290,10 +296,11 @@ class Sheet: }) 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="inbox", + 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) diff --git a/backend/inbox/views.py b/backend/inbox/views.py index 874c9e6..3aed4a0 100644 --- a/backend/inbox/views.py +++ b/backend/inbox/views.py @@ -92,21 +92,8 @@ class Email: async def triage_round(self,message_ids): - """Fetch and classify a whole /email/fetch page, bounded by a semaphore. - - Returns {message_id: decision}. The caller replays the page in upstream order, - so pending_match_ids and pending_confirmation_emails keep the exact sequence - they have today. - - Only the upstream GET and the OpenAI call run concurrently, and nothing inside - the gather touches self.session — Depends(get_session) yields ONE AsyncSession, - which cannot be shared across tasks. All DB work stays in the serial replay. - - Two pre-filters run first and cost no tokens: a message already in - inbox_messages was judged an application once, and a message already in - inbox_message_triage has a stored verdict to replay. That is what makes a - repeated fetch free. - """ + """Fetch and classify a whole /email/fetch page, bounded by a semaphore.""" + ids=[str(m) for m in message_ids or [] if m] decisions={} if not ids: diff --git a/backend/main.py b/backend/main.py index 73916e2..5cda791 100644 --- a/backend/main.py +++ b/backend/main.py @@ -22,6 +22,7 @@ from search.app import router as search_router from interview.app import router as interview_router from talent.app import router as talent_router from candidate_forms.app import router as candidate_forms_router +from g_sheet.app import router as g_sheet_router logging.basicConfig(level=logging.INFO,format="%(levelname)-8s %(name)s: %(message)s") logger=logging.getLogger("main") @@ -32,12 +33,14 @@ async def lifespan(app): async with db_lifespan(app): broker_ready=False cv_broker_ready=False + sheet_broker_ready=False llm_ready=False agent_ready=False close_llm=None close_agent=None broker=None cv_broker=None + sheet_broker=None try: from taskiq_management.broker_setup import broker as _broker broker=_broker @@ -52,6 +55,13 @@ async def lifespan(app): cv_broker_ready=True except Exception as exc: logger.warning("taskiq cv broker startup skipped: %s",exc) + try: + from taskiq_management.g_sheet_broker_setup import sheet_broker as _sheet_broker + sheet_broker=_sheet_broker + await sheet_broker.startup() + sheet_broker_ready=True + except Exception as exc: + logger.warning("taskiq sheet broker startup skipped: %s",exc) try: from llm_setup import init_llm,close_llm as _close_llm from agent.agent_setup import init_agent,close_agent as _close_agent @@ -77,6 +87,8 @@ async def lifespan(app): logger.warning("classifier close skipped: %s",exc) if llm_ready and close_llm is not None: await close_llm() + if sheet_broker_ready and sheet_broker is not None: + await sheet_broker.shutdown() if cv_broker_ready and cv_broker is not None: await cv_broker.shutdown() if broker_ready and broker is not None: @@ -116,3 +128,4 @@ app.include_router(search_router) app.include_router(interview_router) app.include_router(talent_router) app.include_router(candidate_forms_router) +app.include_router(g_sheet_router) diff --git a/backend/requirements.txt b/backend/requirements.txt index fc1aad3..3cec005 100644 --- a/backend/requirements.txt +++ b/backend/requirements.txt @@ -45,3 +45,8 @@ langgraph==1.2.10 # StateGraph agent framework in agent/agent_setup.py # pip install -e .. # Its dependencies are already satisfied by the pins above. openpyxl==3.1.5 + +# --- Google Sheets (g_sheet/) ---------------------------------------------- +google-api-python-client==2.198.0 # Sheets v4 client in g_sheet/plugins.py +google-auth==2.56.3 # ADC + refresh in g_sheet/plugins.py +google-auth-httplib2==0.4.1 # transport used by googleapiclient diff --git a/backend/taskiq_management/g_sheet_broker_setup.py b/backend/taskiq_management/g_sheet_broker_setup.py new file mode 100644 index 0000000..5ccd9b2 --- /dev/null +++ b/backend/taskiq_management/g_sheet_broker_setup.py @@ -0,0 +1,52 @@ +"""Taskiq Google Sheet import broker — isolated Redis stream so sheet imports +never sit behind inbox sync / Outlook / CV work. + +Worker: taskiq worker taskiq_management.g_sheet_broker_setup:sheet_broker g_sheet.tasks +""" + +from __future__ import annotations + +import os + +from dotenv import load_dotenv +from taskiq.middlewares import SmartRetryMiddleware +from taskiq_redis import ( + ListRedisScheduleSource, + RedisAsyncResultBackend, + RedisStreamBroker, +) + +from taskiq_management.broker_setup import MAX_RETRIES,RETRY_DELAY +from taskiq_management.middleware import DeadLetterMiddleware + +load_dotenv() + +REDIS_URL=os.getenv("REDIS_URL","redis://localhost:6379/0") +SHEET_QUEUE_NAME=os.getenv("TASKIQ_SHEET_QUEUE_NAME","sheet_import") + +result_backend=RedisAsyncResultBackend(redis_url=REDIS_URL) +sheet_schedule_source=ListRedisScheduleSource( + url=REDIS_URL,prefix="taskiq:schedule:sheet", +) + +sheet_broker=( + RedisStreamBroker( + url=REDIS_URL, + queue_name=SHEET_QUEUE_NAME, + consumer_group_name=os.getenv("TASKIQ_CONSUMER_GROUP","taskiq"), + idle_timeout=int(os.getenv("TASKIQ_IDLE_TIMEOUT_MS","600000")), + ) + .with_result_backend(result_backend) + .with_middlewares( + DeadLetterMiddleware(redis_url=REDIS_URL), + SmartRetryMiddleware( + default_retry_count=MAX_RETRIES, + default_retry_label=True, + default_delay=RETRY_DELAY, + use_jitter=True, + use_delay_exponent=True, + max_delay_exponent=float(os.getenv("TASKIQ_MAX_DELAY","120")), + schedule_source=sheet_schedule_source, + ), + ) +) diff --git a/backend/tests/test_g_sheet_plugins.py b/backend/tests/test_g_sheet_plugins.py new file mode 100644 index 0000000..a8195a0 --- /dev/null +++ b/backend/tests/test_g_sheet_plugins.py @@ -0,0 +1,290 @@ +"""Unit tests for g_sheet/plugins.py — mapping against the real 32-column sheet. + +No DB, no network. Pins header aliases, parsers, and the Bilal sample row. +""" + +from __future__ import annotations + +from datetime import datetime, timezone + +from g_sheet import plugins +from g_sheet.enums import FORM_DATA_FIELDS, HEADER_ALIASES, FormDataField +from g_sheet.models import FormData + + +# Real header row from "PK - Recruitment Tracking Sheet" (32 columns). +HEADERS = [ + "um", + "Year", + "Month", + "Date", + "Time of Entry", + "Screened By", + "Candidate Name", + "HR comments", + "Candidate Number", + "Candidate Email", + "Profile Link", + "Area of Expertise", + "Requisition Number", + "Position Suitable For", + "Source of Application", + "Age", + "Marital Status", + "Education", + "University of Graduation ", + "Experience", + "Experience Details", + "Area of Residence", + "Communication Skills\n(01 to 10)", + "Preferred Timings?", + "Availability to work in the H.O", + "Current Company", + "Reason for Leaving", + "How soon can you join", + "Current Salary", + "Expected Salary", + "Pros", + "Cons", +] + +EXPECTED_FIELDS = ( + FormDataField.SERIAL_NO, + FormDataField.ENTRY_YEAR, + FormDataField.ENTRY_MONTH, + FormDataField.ENTRY_DATE, + FormDataField.ENTRY_TIME, + FormDataField.SCREENED_BY, + FormDataField.NAME, + FormDataField.HR_COMMENTS, + FormDataField.CANDIDATE_NUMBER, + FormDataField.CANDIDATE_EMAIL, + FormDataField.PROFILE_LINK, + FormDataField.AREA_OF_EXPERTISE, + FormDataField.REQUISITION_NUMBER, + FormDataField.POSITION_SUITABLE_FOR, + FormDataField.SOURCE_OF_APPLICATION, + FormDataField.AGE, + FormDataField.MARITAL_STATUS, + FormDataField.DEGREE, + FormDataField.UNIVERSITY, + FormDataField.EXPERIENCE, + FormDataField.EXPERIENCE_DETAILS, + FormDataField.AREA_OF_RESIDENCE, + FormDataField.COMMUNICATION_SKILLS, + FormDataField.PREFERRED_TIMINGS, + FormDataField.HO_AVAILABILITY, + FormDataField.CURRENT_COMPANY, + FormDataField.REASON_FOR_LEAVING, + FormDataField.NOTICE_PERIOD, + FormDataField.CURRENT_SALARY, + FormDataField.EXPECTED_SALARY, + FormDataField.PROS, + FormDataField.CONS, +) + +SAMPLE_RECORD = { + "um": "1", + "Year": "2021", + "Month": "June", + "Date": "23-Jun-2021", + "Time of Entry": "10:30 AM", + "Screened By": "Sara", + "Candidate Name": "Muhammad Bilal Khan", + "HR comments": "Good profile", + "Candidate Number": "0303-2892503", + "Candidate Email": "bilal_kf@yahoo.com", + "Profile Link": "https://example.com/bilal", + "Area of Expertise": "Software Development", + "Requisition Number": "REQ-1", + "Position Suitable For": "Backend Engineer", + "Source of Application": "Referral", + "Age": "33", + "Marital Status": "Married", + "Education": "BS CS", + "University of Graduation ": "NUST", + "Experience": "11 Years", + "Experience Details": "Software development related experience", + "Area of Residence": "Islamabad", + "Communication Skills\n(01 to 10)": "7", + "Preferred Timings?": "Morning", + "Availability to work in the H.O": "Yes", + "Current Company": "Acme", + "Reason for Leaving": "Growth", + "How soon can you join": "1 month", + "Current Salary": "110k", + "Expected Salary": "150k", + "Pros": "Strong backend", + "Cons": "Limited cloud", +} + + +def test_every_real_header_maps_and_none_are_unmapped(): + assert len(HEADERS) == 32 + for header, field in zip(HEADERS, EXPECTED_FIELDS): + assert plugins.match_field(header) == field, header + assert plugins.collect_unmapped_headers(HEADERS) == [] + + +def test_experience_details_does_not_overwrite_experience(): + mapped = plugins.map_record_to_form_data("tab", SAMPLE_RECORD, HEADERS, 2) + assert mapped["experience"] == "11 Years" + assert mapped["experience_details"] == "Software development related experience" + + +def test_candidate_number_and_email_land_in_own_columns(): + mapped = plugins.map_record_to_form_data("tab", SAMPLE_RECORD, HEADERS, 2) + assert mapped["candidate_number"] == "0303-2892503" + assert mapped["candidate_email"] == "bilal_kf@yahoo.com" + assert mapped["name"] == "Muhammad Bilal Khan" + + +def test_date_maps_to_entry_date_not_interview_round(): + mapped = plugins.map_record_to_form_data("tab", SAMPLE_RECORD, HEADERS, 2) + assert mapped["entry_date"] == datetime(2021, 6, 23, tzinfo=timezone.utc) + assert "interview_date" not in mapped + + +def test_marital_status_only_lands_in_marital_status(): + mapped = plugins.map_record_to_form_data("tab", SAMPLE_RECORD, HEADERS, 2) + assert mapped["marital_status"] == "Married" + assert "family_details" not in mapped + assert "interview_status" not in mapped + + +def test_canonical_header_strips_parens_question_and_trailing_space(): + assert plugins.canonical_header("Communication Skills\n(01 to 10)") == "communication skills" + assert plugins.canonical_header("Preferred Timings?") == "preferred timings" + assert plugins.canonical_header("University of Graduation ") == "university of graduation" + + +def test_parse_salary_and_score(): + assert plugins.parse_salary("110k") == 110000 + assert plugins.parse_salary("60k-70k") == 60000 + assert plugins.parse_salary("Negotiable") is None + assert plugins.parse_score("7") == 7 + assert plugins.parse_score("15") is None + + +def test_rows_to_indexed_records_keeps_true_sheet_row_across_blank(): + rows = [ + ["Name", "Age"], + ["Ada", "30"], + ["", ""], + ["Bob", "40"], + ] + indexed = plugins.rows_to_indexed_records(rows) + assert indexed == [ + (2, {"Name": "Ada", "Age": "30"}), + (4, {"Name": "Bob", "Age": "40"}), + ] + # read API contract still drops blanks without exposing indices + assert plugins.rows_to_records(rows) == [ + {"Name": "Ada", "Age": "30"}, + {"Name": "Bob", "Age": "40"}, + ] + + +def test_no_alias_string_appears_under_two_fields(): + seen: dict[str, FormDataField] = {} + for field, aliases in HEADER_ALIASES.items(): + for alias in aliases: + assert alias not in seen, f"{alias!r} under {seen[alias]} and {field}" + seen[alias] = field + + +def test_form_data_fields_match_model(): + assert set(FORM_DATA_FIELDS) == set(FormData.model_fields) + + +def test_sample_row_end_to_end(): + mapped = plugins.map_record_to_form_data( + "PK - Recruitment Tracking Sheet", SAMPLE_RECORD, HEADERS, 2, + ) + assert mapped["experience"] == "11 Years" + assert mapped["experience_details"] == "Software development related experience" + assert mapped["candidate_email"] == "bilal_kf@yahoo.com" + assert mapped["age"] == 33 + assert mapped["current_salary_value"] == 110000 + assert mapped["communication_skills"] == 7 + assert mapped["pros"] == "Strong backend" + assert mapped["cons"] == "Limited cloud" + assert mapped["row_number"] == 2 + assert mapped["sheet"] == "PK - Recruitment Tracking Sheet" + assert mapped.get("job_post_id") is None + + +# Google Form Responses tab — headers differ from the PK screening sheet. +FORM_RESPONSE_HEADERS = [ + "Timestamp", + "Email", + "Full Name", + "Gender", + "Marital Status", + "University", + "Educational Degree", + "Year of Graduation", + "Position Applied For", + "LinkedIn Profile Link", + "Drop your updated resume", + "How soon can you join us?", + "Phone number (03XX-XXXXXXX)", + "Are you willing to relocate?", + "Residing City", + "Residing Country", + "Area of Interest", + "Recruiter", + "Where did you hear about the position you're applying for?", +] + +FORM_RESPONSE_RECORD = { + "Timestamp": "7/2/2026 17:50:19", + "Email": "nusratazra@gmail.com", + "Full Name": "Nusrat Azra", + "Gender": "Female", + "Marital Status": "Single", + "University": "Karachi University", + "Educational Degree": "Masters", + "Year of Graduation": "12/2/2007", + "Position Applied For": "Executive Secretary", + "LinkedIn Profile Link": "https://www.linkedin.com/in/nusrat-azra-executive-manager-to-c-suite-3bb86719/", + "Drop your updated resume": "https://drive.google.com/open?id=1l5HOY5R6KiL6270A_sV56FkdX6EIGfI4", + "How soon can you join us?": "1 - 2 weeks", + "Phone number (03XX-XXXXXXX)": "03312464228", + "Are you willing to relocate?": "Yes", + "Residing City": "Karachi", + "Residing Country": "Pakistan", + "Area of Interest": "", + "Recruiter": "", + "Where did you hear about the position you're applying for?": "Indeed", +} + + +def test_form_response_headers_map_to_typed_columns(): + mapped = plugins.map_record_to_form_data( + "Form Responses - Candidate Database Sheet 2026", + FORM_RESPONSE_RECORD, + FORM_RESPONSE_HEADERS, + 9, + ) + assert mapped["name"] == "Nusrat Azra" + assert mapped["candidate_email"] == "nusratazra@gmail.com" + assert mapped["candidate_number"] == "03312464228" + assert mapped["degree"] == "Masters" + assert mapped["university"] == "Karachi University" + assert mapped["position_suitable_for"] == "Executive Secretary" + assert mapped["notice_period"] == "1 - 2 weeks" + assert mapped["source_of_application"] == "Indeed" + assert mapped["ho_availability"] == "Yes" + assert mapped["marital_status"] == "Single" + assert mapped["area_of_residence"] == "Karachi" # first matching residence header wins + assert mapped["profile_link"] == ( + "https://www.linkedin.com/in/nusrat-azra-executive-manager-to-c-suite-3bb86719/" + ) + assert mapped["entry_year"] == "12/2/2007" + assert mapped["entry_date"] is not None + assert mapped["entry_time"] == "17:50:19" + # Timestamp must never be used as the candidate name. + assert mapped["name"] != "7/2/2026 17:50:19" + assert mapped["job_post_id"] is None + assert mapped["raw_record"]["Full Name"] == "Nusrat Azra" diff --git a/docker-compose.dev.yml b/docker-compose.dev.yml index 19e319c..7a0b992 100644 --- a/docker-compose.dev.yml +++ b/docker-compose.dev.yml @@ -72,3 +72,8 @@ services: - ./backend:/app - ./app:/app/app - ${ATTACHMENTS_DIR:-./backend/inbox/decoded_attachments}:/app/inbox/decoded_attachments + + taskiq-sheet-worker: + volumes: + - ./backend:/app + - ./app:/app/app diff --git a/docker-compose.yml b/docker-compose.yml index b13b9dd..d1d34fc 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -285,6 +285,24 @@ services: <<: *backend-env TASKIQ_CV_QUEUE_NAME: cv_upload + # Dedicated stream: Google Sheet → FormData import must not block inbox/CV/mailbox. + taskiq-sheet-worker: + <<: *backend-service + container_name: hrms-taskiq-sheet-worker + command: + [ + "taskiq", + "worker", + "taskiq_management.g_sheet_broker_setup:sheet_broker", + "g_sheet.tasks", + "--workers", + "1", + ] + environment: + <<: *backend-env + TASKIQ_SHEET_QUEUE_NAME: sheet_import + TASKIQ_WORKER_NAME: sheet-worker-01 + # Dedicated stream: Outlook pull/triage must not block match/ATS or CV uploads. taskiq-mailbox-sync-worker: <<: *backend-service