From ea6f67786c515f6c5751094d3398ba9e6343e60d Mon Sep 17 00:00:00 2001 From: "ahmed.mujtaba" Date: Mon, 24 Aug 2026 15:04:45 +0500 Subject: [PATCH] . --- backend/g_sheet/app.py | 264 ++++++++++++++++ backend/g_sheet/enums.py | 279 +++++++++++++++++ backend/g_sheet/models.py | 210 +++++++++++++ backend/g_sheet/plugins.py | 536 +++++++++++++++++++++++++++++++++ backend/g_sheet/serializers.py | 188 ++++++++++++ backend/g_sheet/tasks.py | 93 ++++++ backend/g_sheet/views.py | 334 ++++++++++++++++++++ 7 files changed, 1904 insertions(+) create mode 100644 backend/g_sheet/app.py create mode 100644 backend/g_sheet/enums.py create mode 100644 backend/g_sheet/models.py create mode 100644 backend/g_sheet/plugins.py create mode 100644 backend/g_sheet/serializers.py create mode 100644 backend/g_sheet/tasks.py create mode 100644 backend/g_sheet/views.py diff --git a/backend/g_sheet/app.py b/backend/g_sheet/app.py new file mode 100644 index 0000000..07f8769 --- /dev/null +++ b/backend/g_sheet/app.py @@ -0,0 +1,264 @@ +from fastapi import APIRouter,Depends,HTTPException,Query +from fastapi.responses import JSONResponse +from pydantic import BaseModel +from sqlalchemy.ext.asyncio import AsyncSession + +from db_setup import get_session +from g_sheet.views import Sheet +from users.permissions import PermissionTag,require_permission +from dotenv import load_dotenv +load_dotenv() + +router = APIRouter() + + +class AppendRowsBody(BaseModel): + rows: list[list[str]] + + +class UpdateRangeBody(BaseModel): + cell_range: str + rows: list[list[str]] + + +class ClearRangeBody(BaseModel): + cell_range: str + + +@router.get("/sheet/health") +async def sheet_health(): + """Liveness for the Sheets integration — credentials + spreadsheet reachability. + + Unauthenticated like GET /health in main.py, and never 500s: an unreachable sheet + comes back as {"status":"error"} so a probe can read the reason. + """ + try: + service=Sheet() + data=await service.health_check() + return JSONResponse(content={"data":data,"total":1,"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + +@router.get("/sheet/metadata") +async def fetch_sheet_metadata( + spreadsheet_id: str | None = Query(None), + current_user: dict = Depends(require_permission(PermissionTag.SETTINGS_VIEW)), +): + try: + service=Sheet(spreadsheet_id=spreadsheet_id) + data=await service.get_metadata() + return JSONResponse(content={"data":data,"total":1,"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + +@router.get("/sheet/tabs") +async def fetch_sheet_tabs( + spreadsheet_id: str | None = Query(None), + current_user: dict = Depends(require_permission(PermissionTag.SETTINGS_VIEW)), +): + try: + service=Sheet(spreadsheet_id=spreadsheet_id) + items=await service.list_tabs() + return JSONResponse(content={"data":items,"total":len(items),"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + +@router.get("/sheet/fetch") +async def fetch_sheet( + tab: str | None = Query(None), + cell_range: str | None = Query(None), + raw: bool = Query(False), + current_user: dict = Depends(require_permission(PermissionTag.SETTINGS_VIEW)), + spreadsheet_id: str | None = Query(None), +): + """No tab -> every tab as records. With a tab -> that tab, header-mapped unless + raw=true, which returns the rows exactly as the sheet stores them.""" + try: + service=Sheet(spreadsheet_id=spreadsheet_id) + if not tab: + data=await service.read_all() + return JSONResponse(content={"data":data["sheets"],"total":data["total"],"status_code":200}) + if raw or cell_range: + data=await service.read_range(tab,cell_range) + return JSONResponse(content={"data":data,"total":data["row_count"],"status_code":200}) + data=await service.read_records(tab) + return JSONResponse(content={"data":data["records"],"total":data["total"],"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + +@router.post("/sheet/import") +async def import_all_sheets( + 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.""" + try: + service=Sheet(session=session) + data=await service.start_import(current_user=current_user,tab=None) + return JSONResponse(content={"data":data,"total":1,"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + +@router.post("/sheet/{tab}/import") +async def import_one_sheet( + tab: str, + current_user: dict = Depends(require_permission(PermissionTag.SETTINGS_EDIT)), + session: AsyncSession = Depends(get_session), +): + """Enqueue a single-tab import. Poll GET /sheet/import/fetch for status.""" + try: + service=Sheet(session=session) + data=await service.start_import(current_user=current_user,tab=tab) + return JSONResponse(content={"data":data,"total":1,"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + +@router.get("/sheet/import/fetch") +async def fetch_sheet_import( + run_id: str | None = Query(None), + current_user: dict = Depends(require_permission(PermissionTag.SETTINGS_VIEW)), + session: AsyncSession = Depends(get_session), +): + try: + service=Sheet(session=session) + data=await service.get_import_run(run_id=run_id) + return JSONResponse(content={"data":data,"total":1,"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + +@router.get("/sheet/form-data/sheets") +async def fetch_form_data_sheets( + current_user: dict = Depends(require_permission(PermissionTag.SETTINGS_VIEW)), + session: AsyncSession = Depends(get_session), +): + try: + service=Sheet(session=session) + data=await service.get_imported_sheets() + return JSONResponse(content={"data":data,"total":data["total"],"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + +@router.get("/sheet/form-data/fetch") +async def fetch_form_data( + sheet: str | None = Query(None), + search: str | None = Query(None), + top: int | None = Query(None), + skip: int = Query(0,ge=0), + current_user: dict = Depends(require_permission(PermissionTag.SETTINGS_VIEW)), + session: AsyncSession = Depends(get_session), +): + try: + service=Sheet(session=session) + items,total=await service.get_form_data(sheet=sheet,search=search,top=top,skip=skip) + return JSONResponse(content={"data":items,"total":total,"status_code":200}) + except HTTPException: + 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( + record_id: int, + current_user: dict = Depends(require_permission(PermissionTag.SETTINGS_VIEW)), + session: AsyncSession = Depends(get_session), +): + try: + service=Sheet(session=session) + data=await service.get_form_data_by_id(record_id) + return JSONResponse(content={"data":data,"total":1,"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + +@router.delete("/sheet/form-data/{tab}/delete") +async def delete_form_data_sheet( + tab: str, + current_user: dict = Depends(require_permission(PermissionTag.SETTINGS_DELETE)), + session: AsyncSession = Depends(get_session), +): + try: + service=Sheet(session=session) + data=await service.delete_sheet_data(tab) + return JSONResponse(content={"data":data,"total":1,"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + +@router.post("/sheet/{tab}/append") +async def append_sheet_rows( + tab: str, + payload: AppendRowsBody, + current_user: dict = Depends(require_permission(PermissionTag.SETTINGS_EDIT)), + spreadsheet_id: str | None = Query(None), +): + try: + service=Sheet(spreadsheet_id=spreadsheet_id) + data=await service.append_rows(tab,payload.rows) + return JSONResponse(content={"data":data,"total":1,"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + +@router.patch("/sheet/{tab}/update") +async def update_sheet_range( + tab: str, + payload: UpdateRangeBody, + current_user: dict = Depends(require_permission(PermissionTag.SETTINGS_EDIT)), + spreadsheet_id: str | None = Query(None), +): + try: + service=Sheet(spreadsheet_id=spreadsheet_id) + data=await service.update_range(tab,payload.cell_range,payload.rows) + return JSONResponse(content={"data":data,"total":1,"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) + + +@router.post("/sheet/{tab}/clear") +async def clear_sheet_range( + tab: str, + payload: ClearRangeBody, + current_user: dict = Depends(require_permission(PermissionTag.SETTINGS_EDIT)), + spreadsheet_id: str | None = Query(None), +): + try: + service=Sheet(spreadsheet_id=spreadsheet_id) + data=await service.clear_range(tab,payload.cell_range) + return JSONResponse(content={"data":data,"total":1,"status_code":200}) + except HTTPException: + raise + except Exception as e: + raise HTTPException(status_code=500,detail=str(e)) diff --git a/backend/g_sheet/enums.py b/backend/g_sheet/enums.py new file mode 100644 index 0000000..0232c91 --- /dev/null +++ b/backend/g_sheet/enums.py @@ -0,0 +1,279 @@ +"""Sheet header aliases, FormData keys, and date/round 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. +""" + +from enum import Enum + + +class FormDataField(str, Enum): + """Canonical FormData column keys (plus title, which stays in JSONB only).""" + + NAME = "name" + DEGREE = "degree" + EXPERIENCE = "experience" + 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" + DEGREE = "degree" + QUALIFICATION = "qualification" + + @classmethod + def has(cls, value) -> bool: + return value in cls._value2member_map_ + + +class ExperienceAlias(str, Enum): + 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_ + + +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, +} + + +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" + + +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) + + +class NotesToken(str, Enum): + """Substrings that classify a header as RoundRole.NOTES.""" + + NOTE = "note" + REMARK = "remark" + COMMENT = "comment" + + @classmethod + def contained_in(cls, text: str) -> bool: + return any(member.value in text for member in cls) + + +# -- 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) + + +# -- Date parsing ------------------------------------------------------------ + +class DateFormat(str, Enum): + """strptime patterns tried in definition order. + + DD/MM before MM/DD: 14/10/20 is ambiguous and DD/MM is the local convention. + """ + + D_MON_Y_DASH = "%d-%b-%Y" + D_MONTH_Y_DASH = "%d-%B-%Y" + D_MON_Y_SPACE = "%d %b %Y" + D_MONTH_Y_SPACE = "%d %B %Y" + DMY_SLASH = "%d/%m/%Y" + DMY_SLASH_SHORT = "%d/%m/%y" + DMY_DASH = "%d-%m-%Y" + DMY_DASH_SHORT = "%d-%m-%y" + ISO = "%Y-%m-%d" + DMY_DOT = "%d.%m.%Y" + DMY_DOT_SHORT = "%d.%m.%y" + MDY_SLASH = "%m/%d/%Y" + MDY_SLASH_SHORT = "%m/%d/%y" + MON_D_Y = "%b %d %Y" + MONTH_D_Y = "%B %d %Y" + D_MON_Y_SHORT = "%d-%b-%y" + D_MON_Y_SPACE_SHORT = "%d %b %y" + D_MON_Y_SLASH = "%d/%b/%Y" + D_MON_Y_SLASH_SHORT = "%d/%b/%y" + + +class DateTimeSeparator(str, Enum): + """Separators that split a date cell into date + time tails.""" + + DASH = " - " + EN_DASH = " – " + EM_DASH = " — " + SLASH_SPACE = "/ " + PIPE = " | " + + +class MonthNormalisation(Enum): + """Sheet month spellings → %b-safe short form. value is (source, short).""" + + SEPTEMBER = ("september", "sep") + SEPT = ("sept", "sep") + JULY = ("july", "jul") + JUNE = ("june", "jun") + APRIL = ("april", "apr") + MARCH = ("march", "mar") + + @property + def source(self) -> str: + return self.value[0] + + @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 new file mode 100644 index 0000000..502e0c8 --- /dev/null +++ b/backend/g_sheet/models.py @@ -0,0 +1,210 @@ +"""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, 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) + + +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: int | None = Field(default=None, primary_key=True) + sheet: str = Field(nullable=False, index=True) + name: str | None = Field(default=None, index=True) + degree: str | None = Field(default=None) + experience: 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) + + 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)) + + @classmethod + def _filters(cls, *, sheet=None, search=None): + filters = [] + if sheet: + filters.append(cls.sheet == sheet) + if search: + pattern = f"%{search}%" + filters.append(or_( + cls.name.ilike(pattern), + cls.degree.ilike(pattern), + cls.experience.ilike(pattern), + cls.interview_by.ilike(pattern), + )) + return filters + + @classmethod + async def get_form_data_by_id(cls, session: AsyncSession, record_id): + try: + rid = int(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): + rows = [cls(**fields) for fields in records] + session.add_all(rows) + if commit: + await session.commit() + return len(rows) + + @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) + created_by: uuid.UUID | None = Field(default=None, foreign_key="users.id") + 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 diff --git a/backend/g_sheet/plugins.py b/backend/g_sheet/plugins.py new file mode 100644 index 0000000..34001b2 --- /dev/null +++ b/backend/g_sheet/plugins.py @@ -0,0 +1,536 @@ +"""Google Sheets helpers — credential loading, retrying API calls, row/record shaping. + +No FastAPI imports here by house rule: this module raises its own SheetsServiceError +family and lets g_sheet/views.py translate that into HTTPException. + +Auth reuses the credentials already on disk (authorized_user ADC + a valid refresh +token). Nothing here launches a browser, runs InstalledAppFlow, or reads stdin. +""" + +from __future__ import annotations + +import logging +import os +import random +import re +import time +from datetime import datetime, timezone +from pathlib import Path + +from dotenv import load_dotenv +from google.auth import default as google_auth_default +from google.auth.transport.requests import Request +from googleapiclient.discovery import build +from googleapiclient.errors import HttpError + +from g_sheet.enums import ( + ConductedByAlias, + DateFormat, + DateTimeSeparator, + FIELD_ALIAS_ENUMS, + FormDataField, + MonthNormalisation, + NotesToken, + RoundByColumn, + RoundDateColumn, + RoundNotesColumn, + RoundOrdinal, + RoundResultColumn, + RoundRole, + RoundStatusColumn, + RoundTimeColumn, +) + +load_dotenv() + +logger=logging.getLogger("g_sheet.plugins") + +# backend/ — GOOGLE_APPLICATION_CREDENTIALS is stored relative to it ("credentials/..."). +ROOT=Path(__file__).resolve().parent.parent + +SCOPES=[ + "https://www.googleapis.com/auth/spreadsheets", + "https://www.googleapis.com/auth/drive", +] + +SPREADSHEET_ID=os.getenv("SPREADSHEET_ID") +SPREADSHEET_NAME=os.getenv("SPREADSHEET_NAME") +SPREADSHEET_URL=os.getenv("SPREADSHEET_URL") +GOOGLE_APPLICATION_CREDENTIALS=os.getenv("GOOGLE_APPLICATION_CREDENTIALS") + +# 429 and 5xx are transient; every other 4xx is a bad request that a retry repeats. +RETRY_ATTEMPTS=3 +RETRY_BASE_DELAY=0.5 +RETRY_MAX_DELAY=8.0 +RETRYABLE_STATUSES={429,500,502,503,504} + + +class SheetsServiceError(Exception): + """Base for every failure this domain raises. Carries an HTTP-ish status code.""" + + status_code=500 + + def __init__(self,message,status_code=None): + super().__init__(message) + self.message=message + if status_code is not None: + self.status_code=status_code + + +class SheetsAuthError(SheetsServiceError): + """Credentials missing, unreadable, or rejected by Google.""" + + status_code=401 + + +class SheetsApiError(SheetsServiceError): + """The Sheets API answered with an error. status_code is Google's own.""" + + status_code=502 + + +def resolve_credentials_path(credentials_path=None): + """Absolute path to the ADC json. Relative values resolve against backend/. + + The service may be imported from any working directory, so a bare + "credentials/application_default_credentials.json" must not depend on cwd. + """ + raw=credentials_path or GOOGLE_APPLICATION_CREDENTIALS + if not raw: + return None + path=Path(raw) + if not path.is_absolute(): + path=ROOT/path + return path + + +def load_credentials(credentials_path=None,scopes=None): + """Build scoped ADC credentials and refresh them once. Never prompts.""" + path=resolve_credentials_path(credentials_path) + if path is not None: + if not path.exists(): + raise SheetsAuthError(f"Google credentials file not found: {path.name}") + os.environ["GOOGLE_APPLICATION_CREDENTIALS"]=str(path) + try: + credentials,_=google_auth_default(scopes=scopes or SCOPES) + credentials.refresh(Request()) + except SheetsServiceError: + raise + except Exception as e: + raise SheetsAuthError(f"Google credential refresh failed: {e}") + return credentials + + +def ensure_fresh(credentials): + """Refresh only when the token has actually gone stale — not on every call.""" + if credentials is None: + raise SheetsAuthError("Google credentials are not initialised") + if credentials.valid and not credentials.expired: + return credentials + try: + credentials.refresh(Request()) + except Exception as e: + raise SheetsAuthError(f"Google credential refresh failed: {e}") + return credentials + + +def build_sheets_client(credentials): + """Sheets v4 client. cache_discovery=False — the file cache warns under threads.""" + try: + return build("sheets","v4",credentials=credentials,cache_discovery=False) + except Exception as e: + raise SheetsApiError(f"Could not build the Sheets client: {e}") + + +def _status_of(error): + status=getattr(getattr(error,"resp",None),"status",None) + if status is None: + status=getattr(error,"status_code",None) + try: + return int(status) + except (TypeError,ValueError): + return None + + +def _reason_of(error): + """Google's message without the response body, so nothing sensitive leaks out.""" + try: + return error._get_reason().strip() + except Exception: + return str(error) + + +def execute(request,description="sheets request"): + """Run a googleapiclient request with jittered exponential backoff. + + Retries 429 and 5xx up to RETRY_ATTEMPTS; every other HttpError raises straight + away as SheetsApiError carrying Google's status code. + """ + delay=RETRY_BASE_DELAY + last_error=None + for attempt in range(1,RETRY_ATTEMPTS+1): + try: + return request.execute() + except HttpError as e: + status=_status_of(e) + reason=_reason_of(e) + last_error=SheetsApiError(f"{description} failed: {reason}",status or 502) + if status not in RETRYABLE_STATUSES or attempt==RETRY_ATTEMPTS: + raise last_error + sleep_for=min(delay,RETRY_MAX_DELAY)+random.uniform(0,RETRY_BASE_DELAY) + logger.warning( + "%s got %s, retry %s/%s in %.2fs", + description,status,attempt,RETRY_ATTEMPTS,sleep_for, + ) + time.sleep(sleep_for) + delay*=2 + except SheetsServiceError: + raise + except Exception as e: + raise SheetsApiError(f"{description} failed: {e}") + raise last_error + + +def quote_tab(tab,cell_range=None): + """A1 target for a tab whose name may contain spaces or quotes.""" + safe=str(tab).replace("'","''") + if cell_range: + return f"'{safe}'!{cell_range}" + return f"'{safe}'" + + +def normalise_headers(header_row): + """First row -> unique, non-empty column keys. + + Blank cells become column_{i}; a repeated header keeps its first spelling and the + later ones get _1, _2 so no key silently overwrites another. + """ + headers=[] + seen={} + for index,raw in enumerate(header_row): + name=str(raw).strip() if raw is not None else "" + if not name: + name=f"column_{index}" + count=seen.get(name,0) + seen[name]=count+1 + headers.append(name if count==0 else f"{name}_{count}") + return headers + + +def rows_to_records(rows): + """Sheet rows -> list of dicts keyed by the header row. + + Sheets truncates trailing empties, so short rows are padded to header width. + Fully blank rows are dropped rather than emitted as all-empty records. + """ + if not rows: + return [] + headers=normalise_headers(rows[0]) + records=[] + for row in 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)0: + date_part=date_part[:time_match.start()].strip(" ,;-") + + date_part=_DAY_ORDINAL_RE.sub(r"\1",date_part) + date_part=_normalise_month_spellings(date_part) + date_part=re.sub(r"\s+"," ",date_part).strip(" ,;") + + for fmt in DateFormat: + try: + return datetime.strptime(date_part,fmt.value).replace(tzinfo=timezone.utc) + except ValueError: + continue + return None + + +def parse_date_time(value): + """(datetime|None, time_string|None) — fills *_time for the cells that carry one.""" + parsed=parse_date(value) + if value is None: + return parsed,None + text=str(value).strip() + match=_TIME_RE.search(text) + time_str=match.group(1).strip() if match else None + return parsed,time_str + + +def parse_age(value): + """(int|None, raw|None) — first digit run if 0 < n < 100, always keep the raw.""" + if value is None: + return None,None + raw=str(value).strip() + if not raw: + return None,None + match=_AGE_RE.search(raw) + if not match: + return None,raw + number=int(match.group()) + if 0 dict: + """spreadsheets.get response -> the spreadsheet header the UI renders.""" + properties = payload.get("properties") or {} + return { + "spreadsheet_id": payload.get("spreadsheetId"), + "title": properties.get("title"), + "locale": properties.get("locale"), + "time_zone": properties.get("timeZone"), + "url": payload.get("spreadsheetUrl"), + "tabs": [serialize_tab(sheet) for sheet in payload.get("sheets") or []], + } + + +def serialize_tab(sheet: dict) -> dict: + """One entry of spreadsheets.get -> tab name plus its grid size.""" + properties = sheet.get("properties") or {} + grid = properties.get("gridProperties") or {} + return { + "title": properties.get("title"), + "sheet_id": properties.get("sheetId"), + "index": properties.get("index"), + "row_count": grid.get("rowCount"), + "column_count": grid.get("columnCount"), + } + + +def serialize_values(tab: str, cell_range: str | None, rows: list[list[str]]) -> dict: + """Raw rows -> the read_range envelope.""" + return { + "tab": tab, + "range": cell_range, + "rows": rows, + "row_count": len(rows), + } + + +def serialize_records(tab: str, records: list[dict]) -> dict: + """Header-mapped rows -> the read_records envelope.""" + return { + "tab": tab, + "records": records, + "total": len(records), + "headers": list(records[0].keys()) if records else [], + } + + +def serialize_append(tab: str, payload: dict) -> dict: + """values.append response -> what was written and where.""" + updates = payload.get("updates") or {} + return { + "tab": tab, + "spreadsheet_id": payload.get("spreadsheetId"), + "updated_range": updates.get("updatedRange"), + "updated_rows": updates.get("updatedRows", 0), + "updated_columns": updates.get("updatedColumns", 0), + "updated_cells": updates.get("updatedCells", 0), + } + + +def serialize_update(tab: str, payload: dict) -> dict: + """values.update response -> the same shape as an append result.""" + return { + "tab": tab, + "spreadsheet_id": payload.get("spreadsheetId"), + "updated_range": payload.get("updatedRange"), + "updated_rows": payload.get("updatedRows", 0), + "updated_columns": payload.get("updatedColumns", 0), + "updated_cells": payload.get("updatedCells", 0), + } + + +def serialize_clear(tab: str, payload: dict) -> dict: + """values.clear response -> the cleared range.""" + return { + "tab": tab, + "spreadsheet_id": payload.get("spreadsheetId"), + "cleared_range": payload.get("clearedRange"), + } + + +def serialize_health(ok: bool, detail: str, tabs: list[str] | None = None) -> dict: + """health_check result. Returned on failure too — this one never raises.""" + return { + "status": "ok" if ok else "error", + "detail": detail, + "tabs": tabs or [], + "tab_count": len(tabs or []), + } + + +def _iso(value): + return value.isoformat() if value is not None else None + + +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), + } + + +def serialize_import(report: dict) -> dict: + """Per-tab import report.""" + return { + "tab": report.get("tab"), + "rows_read": report.get("rows_read", 0), + "inserted": report.get("inserted", 0), + "deleted": report.get("deleted", 0), + "dates_parsed": report.get("dates_parsed", 0), + "dates_unparsed": report.get("dates_unparsed", 0), + "ages_parsed": report.get("ages_parsed", 0), + "unmapped_headers": report.get("unmapped_headers") or [], + "error": report.get("error"), + } + + +def serialize_import_all(reports: list[dict]) -> dict: + """Aggregate of per-tab reports from import_all.""" + ok=[r for r in reports if not r.get("error")] + failed=[r for r in reports if r.get("error")] + return { + "tabs": len(reports), + "succeeded": len(ok), + "failed": len(failed), + "inserted": sum(r.get("inserted", 0) for r in ok), + "deleted": sum(r.get("deleted", 0) for r in ok), + "reports": [serialize_import(r) for r in reports], + } + + +def serialize_sheet_summary(sheets: list[str]) -> dict: + return {"sheets": sheets, "total": len(sheets)} + + +def serialize_import_run(row) -> dict: + return { + "id": str(row.id), + "status": row.status, + "task_id": row.task_id, + "created_by": str(row.created_by) if row.created_by else None, + "tab": row.tab, + "report": row.report, + "error": row.error, + "created_at": _iso(row.created_at), + "started_at": _iso(row.started_at), + "finished_at": _iso(row.finished_at), + } diff --git a/backend/g_sheet/tasks.py b/backend/g_sheet/tasks.py new file mode 100644 index 0000000..488b2d5 --- /dev/null +++ b/backend/g_sheet/tasks.py @@ -0,0 +1,93 @@ +"""Google Sheet → FormData import Taskiq tasks (shared inbox worker stream).""" + +from __future__ import annotations + +import logging +import os +from datetime import datetime,timezone + +import redis.asyncio as redis +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.middleware import PermanentTaskError + +load_dotenv() + +logger=logging.getLogger("g_sheet.tasks") +REDIS_URL=os.getenv("REDIS_URL","redis://localhost:6379/0") +_LOCK_KEY="g_sheet:import:lock" +_LOCK_TTL=3600 + + +async def _fail(run_id:str,error:str) -> dict: + async with session_scope() as session: + await SheetImportRun.update_run(session,run_id,{ + "status":"failed", + "error":error, + "finished_at":datetime.now(timezone.utc), + }) + return {"status":"failed","error":error} + + +@broker.task( + task_name="g_sheet.import_sheets", + retry_on_error=True, + max_retries=MAX_RETRIES, + delay=RETRY_DELAY, +) +async def import_sheets(run_id:str) -> dict: + if not run_id or not str(run_id).strip(): + raise PermanentTaskError("run_id is required") + run_id=str(run_id).strip() + + client=redis.from_url(REDIS_URL,decode_responses=True) + try: + acquired=await client.set(_LOCK_KEY,run_id,nx=True,ex=_LOCK_TTL) + if not acquired: + return await _fail(run_id,"another sheet import is already running") + + try: + async with session_scope() as session: + row=await SheetImportRun.get_by_id(session,run_id) + if not row: + raise PermanentTaskError(f"import run {run_id} not found") + await SheetImportRun.update_run(session,run_id,{ + "status":"running", + "started_at":datetime.now(timezone.utc), + "error":None, + }) + tab=row.tab + + async with session_scope() as session: + service=Sheet(session=session) + try: + if tab: + report=await service.import_sheet(tab) + else: + report=await service.import_all() + except Exception as e: + logger.exception("sheet import failed for run %s",run_id) + # Bad tab names and permanent Sheets 4xx — do not burn retries. + from fastapi import HTTPException + if isinstance(e,HTTPException) and e.status_code in (400,404,422): + await _fail(run_id,str(e.detail)) + raise PermanentTaskError(str(e.detail)) from e + return await _fail(run_id,str(e)) + + await SheetImportRun.update_run(session,run_id,{ + "status":"completed", + "report":report, + "error":None, + "finished_at":datetime.now(timezone.utc), + }) + return {"status":"completed","report":report} + finally: + current=await client.get(_LOCK_KEY) + if current==run_id: + await client.delete(_LOCK_KEY) + finally: + await client.aclose() diff --git a/backend/g_sheet/views.py b/backend/g_sheet/views.py new file mode 100644 index 0000000..96bf700 --- /dev/null +++ b/backend/g_sheet/views.py @@ -0,0 +1,334 @@ +"""Google Sheets service — business logic for the g_sheet domain. + +The Google client is blocking, so every call goes through asyncio.to_thread rather +than stalling the event loop. Client construction is lazy and guarded by a lock so +concurrent requests build it exactly once. +""" + +import asyncio +import logging +import threading +from datetime import datetime,timezone + +from fastapi import HTTPException + +from g_sheet.plugins import ( + SCOPES, + SPREADSHEET_ID, + SPREADSHEET_NAME, + SPREADSHEET_URL, + SheetsServiceError, + build_sheets_client, + ensure_fresh, + execute, + import_row_stats, + load_credentials, + map_record_to_form_data, + quote_tab, + rows_to_records, + stringify_rows, +) +from g_sheet.models import FormData,SheetImportRun +from g_sheet.serializers import ( + serialize_append, + serialize_clear, + serialize_form_data, + serialize_health, + serialize_import, + serialize_import_all, + serialize_import_run, + serialize_metadata, + serialize_records, + serialize_sheet_summary, + serialize_update, + serialize_values, +) + +logger=logging.getLogger("g_sheet.views") + + +class Sheet: + def __init__(self,session=None,spreadsheet_id=None,credentials_path=None,scopes=None): + self.session=session + self.spreadsheet_id=spreadsheet_id or SPREADSHEET_ID + self.spreadsheet_name=SPREADSHEET_NAME + self.spreadsheet_url=SPREADSHEET_URL + self.credentials_path=credentials_path + self.scopes=scopes or SCOPES + self.credentials=None + self.client=None + self._lock=threading.Lock() + + def _require_session(self): + if self.session is None: + raise HTTPException(status_code=500,detail="Database session is required") + return self.session + + # -- client ------------------------------------------------------------ + + def _connect(self): + """Build credentials + client once, then keep refreshing the same token. + + Double-checked under the lock: two requests racing here must not each build + their own client. + """ + if self.client is not None: + return ensure_fresh(self.credentials) and self.client + with self._lock: + if self.client is None: + self.credentials=load_credentials(self.credentials_path,self.scopes) + self.client=build_sheets_client(self.credentials) + else: + ensure_fresh(self.credentials) + return self.client + + async def _values(self): + if not self.spreadsheet_id: + raise HTTPException(status_code=500,detail="SPREADSHEET_ID is not configured") + client=await asyncio.to_thread(self._connect) + return client.spreadsheets().values() + + async def _spreadsheets(self): + if not self.spreadsheet_id: + raise HTTPException(status_code=500,detail="SPREADSHEET_ID is not configured") + client=await asyncio.to_thread(self._connect) + return client.spreadsheets() + + # -- reads ------------------------------------------------------------- + + async def get_metadata(self): + """Spreadsheet title, id, url and every tab with its row/column counts.""" + try: + spreadsheets=await self._spreadsheets() + request=spreadsheets.get(spreadsheetId=self.spreadsheet_id,fields=( + "spreadsheetId,spreadsheetUrl,properties(title,locale,timeZone)," + "sheets(properties(sheetId,title,index,gridProperties(rowCount,columnCount)))" + )) + payload=await asyncio.to_thread(execute,request,"spreadsheet metadata") + return serialize_metadata(payload) + except SheetsServiceError as e: + raise HTTPException(status_code=e.status_code,detail=e.message) + + async def list_tabs(self): + """Tab titles in sheet order.""" + metadata=await self.get_metadata() + return [tab["title"] for tab in metadata["tabs"] if tab.get("title")] + + async def read_range(self,tab,cell_range=None): + """Raw rows for a tab, or for a sub-range of it when cell_range is given.""" + try: + values=await self._values() + target=quote_tab(tab,cell_range) + request=values.get(spreadsheetId=self.spreadsheet_id,range=target) + payload=await asyncio.to_thread(execute,request,f"read {target}") + rows=stringify_rows(payload.get("values")) + return serialize_values(tab,cell_range,rows) + except SheetsServiceError as e: + raise HTTPException(status_code=e.status_code,detail=e.message) + + async def read_records(self,tab): + """Rows keyed by the first row. Blank rows are skipped, short rows padded.""" + data=await self.read_range(tab) + return serialize_records(tab,rows_to_records(data["rows"])) + + async def read_all(self): + """Every tab as records, keyed by tab name.""" + tabs=await self.list_tabs() + sheets={} + for tab in tabs: + data=await self.read_records(tab) + sheets[tab]=data["records"] + return {"sheets":sheets,"tabs":tabs,"total":len(tabs)} + + # -- writes ------------------------------------------------------------ + + async def append_rows(self,tab,rows): + """Append rows below the tab's current content.""" + if not rows: + raise HTTPException(status_code=422,detail="rows must not be empty") + try: + values=await self._values() + target=quote_tab(tab) + request=values.append( + spreadsheetId=self.spreadsheet_id, + range=target, + valueInputOption="USER_ENTERED", + insertDataOption="INSERT_ROWS", + body={"values":rows}, + ) + payload=await asyncio.to_thread(execute,request,f"append to {target}") + return serialize_append(tab,payload) + except SheetsServiceError as e: + raise HTTPException(status_code=e.status_code,detail=e.message) + + async def update_range(self,tab,cell_range,rows): + """Overwrite an explicit A1 range with rows.""" + if not cell_range: + raise HTTPException(status_code=422,detail="cell_range is required") + if not rows: + raise HTTPException(status_code=422,detail="rows must not be empty") + try: + values=await self._values() + target=quote_tab(tab,cell_range) + request=values.update( + spreadsheetId=self.spreadsheet_id, + range=target, + valueInputOption="USER_ENTERED", + body={"values":rows}, + ) + payload=await asyncio.to_thread(execute,request,f"update {target}") + return serialize_update(tab,payload) + except SheetsServiceError as e: + raise HTTPException(status_code=e.status_code,detail=e.message) + + async def clear_range(self,tab,cell_range): + """Clear the values in an explicit A1 range, leaving formatting intact.""" + if not cell_range: + raise HTTPException(status_code=422,detail="cell_range is required") + try: + values=await self._values() + target=quote_tab(tab,cell_range) + request=values.clear(spreadsheetId=self.spreadsheet_id,range=target,body={}) + payload=await asyncio.to_thread(execute,request,f"clear {target}") + return serialize_clear(tab,payload) + except SheetsServiceError as e: + raise HTTPException(status_code=e.status_code,detail=e.message) + + # -- FormData import / query ------------------------------------------- + + async def import_sheet(self,tab): + """Read one tab from Google Sheets and replace its FormData rows.""" + session=self._require_session() + if not tab or not str(tab).strip(): + raise HTTPException(status_code=422,detail="tab is required") + tab=str(tab).strip() + data=await self.read_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)) + result=await FormData.replace_sheet(session,tab,mapped) + stats=import_row_stats(mapped,headers) + return serialize_import({ + "tab":tab, + "rows_read":len(records), + "inserted":result["inserted"], + "deleted":result["deleted"], + **stats, + }) + + async def import_all(self): + """Import every tab sequentially; one tab failure does not abort the rest.""" + self._require_session() + tabs=await self.list_tabs() + reports=[] + for tab in tabs: + try: + report=await self.import_sheet(tab) + reports.append(report) + except HTTPException as e: + logger.warning("import_all tab %s failed: %s",tab,e.detail) + reports.append(serialize_import({ + "tab":tab,"rows_read":0,"inserted":0,"deleted":0, + "error":str(e.detail), + })) + except Exception as e: + logger.exception("import_all tab %s failed",tab) + reports.append(serialize_import({ + "tab":tab,"rows_read":0,"inserted":0,"deleted":0, + "error":str(e), + })) + return serialize_import_all(reports) + + async def get_form_data(self,sheet=None,search=None,top=None,skip=None): + session=self._require_session() + rows=await FormData.fetch_form_data( + session,sheet=sheet,search=search,top=top,skip=skip, + ) + total=await FormData.count_form_data(session,sheet=sheet,search=search) + return [serialize_form_data(row) for row in rows],total + + async def get_form_data_by_id(self,record_id): + session=self._require_session() + row=await FormData.get_form_data_by_id(session,record_id) + if not row: + raise HTTPException(status_code=404,detail="Form data not found") + return serialize_form_data(row) + + async def get_imported_sheets(self): + session=self._require_session() + sheets=await FormData.get_sheet_names(session) + return serialize_sheet_summary(sheets) + + async def delete_sheet_data(self,tab): + session=self._require_session() + if not tab or not str(tab).strip(): + raise HTTPException(status_code=422,detail="tab is required") + deleted=await FormData.delete_by_sheet(session,str(tab).strip()) + return {"tab":str(tab).strip(),"deleted":deleted} + + async def start_import(self,current_user=None,tab=None): + """Enqueue a sheet import on the shared Taskiq worker; return the run row. + + If a queued/running import already exists, return it instead of stacking another. + """ + session=self._require_session() + active=await SheetImportRun.get_active(session) + if active: + return serialize_import_run(active) + + created_by=None + if isinstance(current_user,dict) and current_user.get("id"): + created_by=SheetImportRun._as_uuid(current_user.get("id")) + + tab_value=str(tab).strip() if tab else None + row=await SheetImportRun.insert_run(session,{ + "status":"queued", + "created_by":created_by, + "tab":tab_value, + }) + + from g_sheet.tasks import import_sheets + task=await import_sheets.kicker().with_labels( + created_at=datetime.now(timezone.utc).isoformat(), + correlation_id=str(row.id), + queue="inbox", + ).kiq(str(row.id)) + row=await SheetImportRun.update_run(session,row.id,{"task_id":task.task_id}) + return serialize_import_run(row) + + async def get_import_run(self,run_id=None): + session=self._require_session() + if run_id: + row=await SheetImportRun.get_by_id(session,run_id) + if not row: + raise HTTPException(status_code=404,detail="Import run not found") + return serialize_import_run(row) + row=await SheetImportRun.get_active(session) + if row: + return serialize_import_run(row) + from sqlmodel import select + result=await session.execute( + select(SheetImportRun).order_by(SheetImportRun.created_at.desc()).limit(1) + ) + row=result.scalars().first() + if not row: + raise HTTPException(status_code=404,detail="No import runs yet") + return serialize_import_run(row) + + # -- health ------------------------------------------------------------ + + async def health_check(self): + """Credentials + sheet reachability as a status dict. Never raises.""" + if not self.spreadsheet_id: + return serialize_health(False,"SPREADSHEET_ID is not configured") + try: + tabs=await self.list_tabs() + return serialize_health(True,"spreadsheet reachable",tabs) + except HTTPException as e: + logger.warning("sheets health check failed: %s",e.detail) + return serialize_health(False,str(e.detail)) + except Exception as e: + logger.warning("sheets health check failed: %s",e) + return serialize_health(False,str(e))