"""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. Hierarchy: Sheet shared config / session └─ SheetClient credentials + spreadsheets client ├─ SheetRead │ ├─ SheetHealth │ └─ SheetImport └─ SheetWrite SheetFormData DB mirror only (no Google client) """ 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, load_credentials, quote_tab, rows_to_indexed_records, 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: """Parent: spreadsheet identity, optional DB session, and shared helpers.""" 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 class SheetClient(Sheet): """Google API client — lazy connect, token refresh, values/spreadsheets handles.""" 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() class SheetRead(SheetClient): """Read-only sheet operations.""" 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)} class SheetWrite(SheetClient): """Mutating sheet operations.""" 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) class SheetHealth(SheetRead): """Credentials + spreadsheet reachability.""" 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)) class SheetImport(SheetRead): """Google Sheet → FormData import + import-run tracking.""" 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_range(tab) rows=data["rows"] if not rows: return serialize_import({"tab":tab,"rows_read":0,"inserted":0,"deleted":0}) indexed=rows_to_indexed_records(rows) mapped=[ FormData.from_sheet_row(tab,row_number,record) for row_number,record in indexed ] result=await FormData.replace_sheet(session,tab,mapped) return serialize_import({ "tab":tab, "rows_read":len(indexed), "inserted":result["inserted"], "deleted":result["deleted"], }) 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 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 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=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) 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) class SheetFormData(Sheet): """FormData DB mirror — query / delete only (no Google client).""" async def _hydrate_job_posts(self,items): """Attach matching job_posts (title == position_applied_for) + assigned_job_post. No AI suggestions — form applicants already name the role. One query for titles on the page, one for any assigned ids. """ if not items: return items from job.job_post.models import JobPosts from job.job_post.serializers import serialize_job_post session=self._require_session() titles=[(item.get("position_applied_for") or "").strip() for item in items] titles=[t for t in titles if t] by_title={} if titles: for post in await JobPosts.get_by_titles(session,titles): key=(post.title or "").strip().lower() payload=serialize_job_post(post) if post.is_deleted or not post.is_active: payload={**payload,"unavailable":True} by_title.setdefault(key,[]).append(payload) assigned_ids=[item.get("job_post_id") for item in items if item.get("job_post_id")] assigned_map={} if assigned_ids: for post in await JobPosts.get_by_ids(session,assigned_ids,active_only=False): assigned_map[str(post.id)]=serialize_job_post(post) for item in items: key=(item.get("position_applied_for") or "").strip().lower() item["job_posts"]=list(by_title.get(key) or []) aid=item.get("job_post_id") item["assigned_job_post"]=assigned_map.get(str(aid)) if aid else None return items async def get_form_data( self,sheet=None,search=None,offset=0,limit=None, processing_state=None,is_duplicate=None, ): session=self._require_session() rows=await FormData.fetch_form_data( session,sheet=sheet,search=search,offset=offset,limit=limit, processing_state=processing_state,is_duplicate=is_duplicate, ) total=await FormData.count_form_data( session,sheet=sheet,search=search, processing_state=processing_state,is_duplicate=is_duplicate, ) items=await self._hydrate_job_posts([serialize_form_data(row) for row in rows]) return items,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") items=await self._hydrate_job_posts([serialize_form_data(row)]) return items[0] async def assign_job_post(self,record_id,job_post_id): """Set or clear form_data.job_post_id (same contract as inbox assign). Setting a job promotes the row into Users + manual_upload_candidate so Candidates / Talent Pool / Pipeline can see it (platform tag: Form). """ session=self._require_session() if job_post_id is not None: from job.job_post.models import JobPosts post=await JobPosts.get_job_post_by_id(session,job_post_id) if not post or post.is_deleted or not post.is_active: raise HTTPException(status_code=404,detail="Job post not found") updated=await FormData.set_job_post(session,record_id,job_post_id) if not updated: raise HTTPException(status_code=404,detail="Form data not found") if job_post_id is not None: await self._promote_to_application(updated) return await self.get_form_data_by_id(record_id) async def set_processing_state(self,record_id,processing_state): allowed=("unread","imported","processed","rejected") if processing_state not in allowed: raise HTTPException(status_code=422,detail=f"processing_state must be one of {', '.join(allowed)}") 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") # Shortlist requires a job — promote (idempotent) then flip the queue label. if processing_state=="processed": if not row.job_post_id: raise HTTPException(status_code=422,detail="Assign a job post before moving to shortlist") await self._promote_to_application(row) updated=await FormData.set_processing_state(session,record_id,processing_state) if not updated: raise HTTPException(status_code=404,detail="Form data not found") return await self.get_form_data_by_id(record_id) async def _promote_to_application(self,form_row): """Create Users + manual_upload_candidate from a form_data row (idempotent). Pipeline / Candidates / Talent Pool all read manual_upload_candidate (or the CANDIDATE user it creates). platform='Form' is the source badge. """ session=self._require_session() from job.candidate.models import Manual_UPLOAD_CANDIDATE from job.history.views import HistoryRecorder from job.history.enums import HistoryEvent if getattr(form_row,"manual_upload_candidate_id",None): existing=await Manual_UPLOAD_CANDIDATE.get_by_id(session,form_row.manual_upload_candidate_id) if existing: if form_row.job_post_id and existing.job_post_id!=form_row.job_post_id: existing.job_post_id=form_row.job_post_id session.add(existing) await session.commit() return existing email=(form_row.candidate_email or "").strip().lower() if not email: raise HTTPException(status_code=422,detail="candidate_email is required to promote this form applicant") if not form_row.job_post_id: raise HTTPException(status_code=422,detail="job_post_id is required to promote this form applicant") existing=await Manual_UPLOAD_CANDIDATE.get_by_email_and_job( session,email,form_row.job_post_id, ) if existing: await FormData.link_manual_upload(session,form_row.id,existing.id) return existing resume=(form_row.resume_link or "").strip() file_name="" if resume: file_name=resume.rsplit("/",1)[-1][:180] or "resume" # Sheet already stores LinkedIn on profile_link — copy it through, do not parse the CV. profile=(form_row.profile_link or "").strip() linkedin_url=None if profile: linkedin_url=profile if profile.lower().startswith("http") else f"https://{profile.lstrip('/')}" row=await Manual_UPLOAD_CANDIDATE.create_manual_upload_candidate(session,{ "candidate_email":email, "candidate_name":(form_row.name or "").strip() or email, "candidate_phone":(form_row.candidate_number or "").strip(), "job_post_id":str(form_row.job_post_id), "current_company":(form_row.current_company or "").strip(), "current_position":(form_row.position_applied_for or "").strip(), "platform":"Form", "apply_via":"form", "experience":(form_row.experience or "").strip(), "status":"PENDING", "file_name":file_name, "file_path":resume, "full_text":"", "linkedin_url":linkedin_url, }) await FormData.link_manual_upload(session,form_row.id,row.id) try: await HistoryRecorder(session).record( HistoryEvent.CANDIDATE_CREATED.value, actor_id=None,user_id=row.user_id, manual_upload_candidate_id=row.id, entity_type="manual_upload_candidate",entity_id=row.id, to_value=row.candidate_email, description="Form",commit=True, ) except Exception: logger.exception("form promote history record failed for %s",form_row.id) return row async def set_duplicate(self,record_id,is_duplicate): if not isinstance(is_duplicate,bool): raise HTTPException(status_code=422,detail="is_duplicate must be a boolean") updated=await FormData.set_duplicate(self._require_session(),record_id,is_duplicate) if not updated: raise HTTPException(status_code=404,detail="Form data not found") return await self.get_form_data_by_id(record_id) async def get_counts(self,sheet=None): return await FormData.count_processing(self._require_session(),sheet=sheet) 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}