import uuid from datetime import datetime, timezone from typing import TYPE_CHECKING, List, Optional from fastapi import HTTPException from sqlalchemy import JSON, DateTime, func, UniqueConstraint from sqlalchemy.exc import IntegrityError from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import selectinload from sqlmodel import Field, Relationship, SQLModel, select if TYPE_CHECKING: from inbox.models import Inbox from users.models import Users from job.job_post.models import JobPosts def _now() -> datetime: return datetime.now(timezone.utc) # Every datetime below is aware (see _now, and the API parses ISO input carrying # an offset), so each column is declared timestamptz. SQLModel maps a bare # `datetime` to TIMESTAMP WITHOUT TIME ZONE, and asyncpg refuses to bind an aware # value to one — "can't subtract offset-naive and offset-aware datetimes" — which # turns every insert here into a 500. Same pairing as job/job_post/models.py. class Manual_UPLOAD_CANDIDATE(SQLModel, table=True): __tablename__ = "manual_upload_candidate" id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True) candidate_email: str = Field(default="") candidate_name: str = Field(default="") candidate_phone: str = Field(default="") job_post_id: uuid.UUID | None = Field(default=None, foreign_key="job_posts.id") full_text: str = Field(default="") current_company: str = Field(default="") # Candidate's role at that company (e.g. "Senior Merchandiser"). Distinct # from job_posts.title — that is the role they applied to, not their own. # Same ALTER-on-a-populated-table reasoning as referral_by below. current_position: str = Field(default="", sa_column_kwargs={"server_default": ""}) # Add Candidate only (POST /candidate/create/candidate). /import never writes this table. apply_via:Optional[str]=Field(nullable=True) user_id: uuid.UUID | None = Field(default=None, foreign_key="users.id") platform: str = Field(default="") created_by: uuid.UUID | None = Field(default=None, foreign_key="users.id") experience: str = Field(default="") # Candidate_application_Status value (PENDING, SCREENING, …). Empty reads as Shortlist. status: str = Field(default="") # Free text, not a users FK: a referrer is often someone outside the system referral_by: str = Field(default="", sa_column_kwargs={"server_default": ""}) file_name: str = Field(default="", sa_column_kwargs={"server_default": ""}) file_path: str = Field(default="", sa_column_kwargs={"server_default": ""}) 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 async def get_all(cls, session: AsyncSession, job_post_id=None, limit=None, offset=0): try: from inbox.models import AtsResults from users.models import Users from job.job_post.models import JobPosts qry=( select( cls.id, cls.candidate_email, cls.user_id, Users.email, cls.job_post_id, Users.name, cls.candidate_phone, JobPosts.title, cls.status, cls.current_company, cls.current_position, cls.experience, cls.created_at, cls.updated_at, AtsResults.id.label("ats_result_id"), AtsResults.overall_score, AtsResults.band, AtsResults.job_post_id.label("ats_job_post_id"), AtsResults.computed_at, AtsResults.candidate_id, AtsResults.user_id.label("ats_user_id"), ) .join(Users,cls.user_id==Users.id) .join(JobPosts,cls.job_post_id==JobPosts.id) .outerjoin( AtsResults, (Users.id==AtsResults.user_id) &(AtsResults.job_post_id==cls.job_post_id) &(AtsResults.is_current==True), # noqa: E712 ) .order_by( AtsResults.overall_score.desc().nulls_last(), cls.created_at.desc(), ) ) if job_post_id: qry=qry.where(cls.job_post_id==job_post_id) if limit is not None: qry=qry.limit(limit).offset(offset) result=await session.execute(qry) rows=[] for row in result.mappings().all(): ats=None if row["ats_result_id"] is not None: ats={ "id":str(row["ats_result_id"]), "overall_score":row["overall_score"], "band":row["band"] or None, "job_post_id":str(row["ats_job_post_id"]) if row["ats_job_post_id"] else None, "computed_at":row["computed_at"].isoformat() if row["computed_at"] else None, "candidate_id":str(row["candidate_id"]) if row["candidate_id"] else None, "user_id":str(row["ats_user_id"]) if row["ats_user_id"] else None, } rows.append({ "id":str(row["id"]) if row["id"] else None, "candidate_email":row["candidate_email"], "user_id":str(row["user_id"]) if row["user_id"] else None, "email":row["email"], "job_post_id":str(row["job_post_id"]) if row["job_post_id"] else None, "name":row["name"], "candidate_phone":row["candidate_phone"], "title":row["title"] or None, "application_status":row["status"] or None, "current_company":row["current_company"] or None, "current_position":row["current_position"] or None, "experience":row["experience"] or None, "created_at":row["created_at"].isoformat() if row["created_at"] else None, "updated_at":row["updated_at"].isoformat() if row["updated_at"] else None, "ats_result":ats, }) return rows except Exception as e: raise HTTPException(status_code=500,detail=str(e)) @classmethod async def count_by_status(cls, session: AsyncSession, job_post_id=None): try: from users.models import Users from job.job_post.models import JobPosts qry=( select(cls.status,func.count()) .select_from(cls) .join(Users,cls.user_id==Users.id) .join(JobPosts,cls.job_post_id==JobPosts.id) .group_by(cls.status) ) if job_post_id: qry=qry.where(cls.job_post_id==job_post_id) result=await session.execute(qry) counts={} for status,n in result.all(): key=(status or "").strip() or "UNKNOWN" counts[key]=counts.get(key,0)+int(n or 0) return counts except Exception as e: raise HTTPException(status_code=500,detail=str(e)) @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 create_manual_upload_candidate(cls, session: AsyncSession, fields: dict): import os from role.models import EnumRoles, Roles from users.models import Users from users.plugins import hash_password email=(fields.get("candidate_email") or "").strip().lower() name=(fields.get("candidate_name") or "").strip() or email default_pw=os.getenv("DEFAULT_CANDIDATE_PASSWORD","Utopia!@#") user=await Users.get_user_by_email(session,email) if not user: role=await Roles.get_role_by_name(session,EnumRoles.CANDIDATE.value) user=await Users.insert_user(session,{ "name":name, "email":email, "role_id":role.id if role else 8, "password":hash_password(default_pw), "is_active":True, "is_deleted":False, }) row=cls( candidate_email=email, candidate_name=name, candidate_phone=(fields.get("candidate_phone") or "").strip(), job_post_id=cls._as_uuid(fields.get("job_post_id")), full_text=fields.get("full_text") or "", current_company=(fields.get("current_company") or "").strip(), current_position=(fields.get("current_position") or "").strip(), apply_via="manual_upload", user_id=user.id, platform=(fields.get("platform") or "").strip(), created_by=cls._as_uuid(fields.get("created_by")), experience=(fields.get("experience") or "").strip(), status=(fields.get("status") or "").strip() or "PENDING", referral_by=(fields.get("referral_by") or "").strip(), file_name=(fields.get("file_name") or "").strip(), file_path=(fields.get("file_path") or "").strip(), ) session.add(row) await session.commit() await session.refresh(row) return row @classmethod async def get_by_user_id(cls, session: AsyncSession, user_id): uid = cls._as_uuid(user_id) if uid is None: return None result = await session.execute( select(cls).where(cls.user_id == uid).order_by(cls.created_at.desc()) ) return result.scalars().first() @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() class Candidates(SQLModel, table=True): __tablename__ = "candidates" __table_args__ = (UniqueConstraint("job_id", "content_sha256"),) id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True) job_id: uuid.UUID = Field(foreign_key="job_posts.id", index=True) source: str = Field(default="upload") # "upload" | "inbox" origin metadata, not an FK filename: str file_path: str | None = Field(default=None) # decoded-attachment path (inbox only) content_sha256: str | None = Field(default=None, index=True) candidate_email: str | None = Field(default=None) candidate_name: str | None = Field(default=None) job_title: str | None = Field(default=None) current_company: str | None = Field(default=None) years_experience: int | None = Field(default=None) match_score: int | None = Field(default=None) # None on failed rows matched_keywords: list[str] = Field(default_factory=list, sa_type=JSON) missing_keywords: list[str] = Field(default_factory=list, sa_type=JSON) summary_critique: str | None = Field(default=None) status: str # "completed" | "failed" error_code: str | None = Field(default=None) error_message: str | None = Field(default=None) model: str | None = Field(default=None) # which OPENAI_MODEL produced the score created_by: uuid.UUID = Field(foreign_key="users.id") created_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) updated_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) @staticmethod def _as_uuid(record_id) -> uuid.UUID | None: try: return uuid.UUID(str(record_id)) except ValueError: return None @classmethod async def get_candidate_by_id(cls, session: AsyncSession, record_id: str): 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_candidates_by_job(cls, session: AsyncSession, job_id: str | None = None): """Leaderboard order: completed by score desc, failures last, ties stable. job_id=None returns the whole pool across jobs (same ordering) for the frontend's unscoped Candidates/Talent Pool views. """ statement = select(cls) if job_id is not None: uid = cls._as_uuid(job_id) if uid is None: return [] statement = statement.where(cls.job_id == uid) statement = statement.order_by( cls.status.asc(), # "completed" < "failed" cls.match_score.desc().nulls_last(), cls.created_at.asc(), ) result = await session.execute(statement) return result.scalars().all() @classmethod async def get_completed_by_email_job(cls, session: AsyncSession, email, job_id): """Newest completed score for this email against this job — keywords when ats_results.candidate_id is NULL (matched-user identity).""" normalized = (email or "").strip().lower() jid = cls._as_uuid(job_id) if not normalized or jid is None: return None result = await session.execute( select(cls) .where( func.lower(cls.candidate_email) == normalized, cls.job_id == jid, cls.status == "completed", ) .order_by(cls.updated_at.desc()) ) return result.scalars().first() @classmethod async def upsert_candidate(cls, session: AsyncSession, fields: dict): existing = None sha = fields.get("content_sha256") if sha: result = await session.execute( select(cls).where(cls.job_id == fields["job_id"], cls.content_sha256 == sha) ) existing = result.scalars().first() if existing is None: row = cls(**fields) session.add(row) try: await session.commit() except IntegrityError: # A concurrent request inserted the same (job_id, sha) first; take over # that row and update it instead. await session.rollback() result = await session.execute( select(cls).where(cls.job_id == fields["job_id"], cls.content_sha256 == sha) ) existing = result.scalars().first() if existing is None: raise else: await session.refresh(row) return row for key, value in fields.items(): setattr(existing, key, value) existing.updated_at = _now() session.add(existing) await session.commit() await session.refresh(existing) return existing class Interviews(SQLModel, table=True): __tablename__ = "interviews" id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True) interview_date: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) interview_time: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) interview_type: str = Field(default="") interview_status: str = Field(default="") inbox_id: int | None = Field(default=None, foreign_key="inbox.id") graph_event_id: str | None = Field(default=None) web_link: str | None = Field(default=None) inbox: Optional["Inbox"] = Relationship( back_populates="interviews", sa_relationship_kwargs={"lazy": "selectin"}, ) @staticmethod def _as_uuid(record_id) -> uuid.UUID | None: try: return uuid.UUID(str(record_id)) except ValueError: return None @classmethod def _with_inbox_message(cls): from inbox.models import Inbox return selectinload(cls.inbox).selectinload(Inbox.messages) @classmethod async def get_interview_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).options(cls._with_inbox_message()).where(cls.id == uid) ) return result.scalars().first() @classmethod async def get_interviews_by_inbox(cls, session: AsyncSession, inbox_id: int): result = await session.execute( select(cls) .options(cls._with_inbox_message()) .where(cls.inbox_id == inbox_id) .order_by(cls.interview_date.desc()) ) return result.scalars().all() @classmethod async def get_interviews_in_range( cls, session: AsyncSession, *, from_date=None, to_date=None, status: str | None = None, top: int | None = None, skip: int = 0, ): statement = select(cls) if from_date is not None: statement = statement.where(cls.interview_date >= from_date) if to_date is not None: statement = statement.where(cls.interview_date < to_date) if status: statement = statement.where(cls.interview_status == status) count_statement = select(func.count()).select_from(statement.subquery()) total = (await session.execute(count_statement)).scalar_one() statement = ( statement.options(cls._with_inbox_message()).order_by(cls.interview_date.asc()) ) if skip: statement = statement.offset(skip) if top is not None: statement = statement.limit(top) result = await session.execute(statement) return list(result.scalars().all()), total @classmethod async def job_titles_by_inbox(cls, session: AsyncSession, inbox_ids) -> dict[int, str]: """Resolve {inbox_id: job_title} for a page of interview rows. Two constraints keep this off get_interviews_in_range: 1. Interviews.inbox is selectin, but Inbox.messages defaults to lazy="select" and raises MissingGreenlet under AsyncSession. Switching that relation to selectin would extra-query the candidate list, pipeline, activity and inbox. 2. Widening the range statement with a join would INNER-join the count subquery and drop interviews whose application has no assigned requisition, changing `total` on Interviews and Calendar. INNER joins are correct here — an unassigned inbox simply produces no dict entry and .get() yields None. No response rows are lost because the rows still come from the untouched range statement. """ from inbox.models import Inbox, Inbox_Messages from job.job_post.models import JobPosts ids = {int(i) for i in (inbox_ids or []) if i is not None} if not ids: return {} result = await session.execute( select(Inbox.id, JobPosts.title) .join(Inbox_Messages, Inbox.message_id == Inbox_Messages.id) .join(JobPosts, Inbox_Messages.assigned_job_post_id == JobPosts.id) .where(Inbox.id.in_(ids)) ) return {int(inbox_id): title for inbox_id, title in result.all()} @classmethod async def insert_interview(cls, session: AsyncSession, fields: dict): row = cls(**fields) session.add(row) await session.commit() return await cls.get_interview_by_id(session, row.id) @classmethod async def update_interview(cls, session: AsyncSession, record_id, fields: dict): row = await cls.get_interview_by_id(session, record_id) if not row: return None for key, value in fields.items(): setattr(row, key, value) session.add(row) await session.commit() await session.refresh(row) return row @classmethod async def set_calendar_event(cls, session: AsyncSession, record_id, event_id, web_link): row = await cls.get_interview_by_id(session, record_id) if not row: return None row.graph_event_id = event_id row.web_link = web_link session.add(row) await session.commit() await session.refresh(row) return row class Notes(SQLModel, table=True): __tablename__ = "notes" id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True) note: str = Field(default="") created_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) updated_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) user_id: uuid.UUID | None = Field(default=None, foreign_key="users.id") created_by: uuid.UUID | None = Field(default=None, foreign_key="users.id") user: Optional["Users"] = Relationship( back_populates="notes", sa_relationship_kwargs={"lazy": "selectin", "foreign_keys": "[Notes.user_id]"}, ) author: Optional["Users"] = Relationship( back_populates="authored_notes", sa_relationship_kwargs={"lazy": "selectin", "foreign_keys": "[Notes.created_by]"}, ) @staticmethod def _as_uuid(record_id) -> uuid.UUID | None: try: return uuid.UUID(str(record_id)) except ValueError: return None @classmethod async def get_note_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_notes_by_user(cls, session: AsyncSession, user_id): uid = cls._as_uuid(user_id) if uid is None: return [] result = await session.execute( select(cls).where(cls.user_id == uid).order_by(cls.created_at.desc()) ) return result.scalars().all() @classmethod async def insert_note(cls, session: AsyncSession, fields: dict): row = cls(**fields) session.add(row) await session.commit() return await cls.get_note_by_id(session, row.id) @classmethod async def update_note(cls, session: AsyncSession, record_id, fields: dict): row = await cls.get_note_by_id(session, record_id) if not row: return None for key, value in fields.items(): setattr(row, key, value) row.updated_at = _now() session.add(row) await session.commit() await session.refresh(row) return row class Activity(SQLModel, table=True): __tablename__ = "activity" id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True) activity_type: str = Field(default="") activity_date: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) activity_time: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) activity_status: str = Field(default="") description: str | None = Field(default=None) inbox_id: int | None = Field(default=None, foreign_key="inbox.id") inbox: Optional["Inbox"] = Relationship( back_populates="activity", sa_relationship_kwargs={"lazy": "selectin"}, ) @staticmethod def _as_uuid(record_id) -> uuid.UUID | None: try: return uuid.UUID(str(record_id)) except ValueError: return None @classmethod async def get_activity_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_activity_by_inbox(cls, session: AsyncSession, inbox_id: int): result = await session.execute( select(cls).where(cls.inbox_id == inbox_id).order_by(cls.activity_date.desc()) ) return result.scalars().all() @classmethod async def get_activity_feed(cls, session: AsyncSession, *, top: int | None = None, skip: int = 0): statement = select(cls) count_statement = select(func.count()).select_from(cls) total = (await session.execute(count_statement)).scalar_one() statement = statement.order_by(cls.activity_date.desc(), cls.activity_time.desc()) if skip: statement = statement.offset(skip) if top is not None: statement = statement.limit(top) result = await session.execute(statement) return list(result.scalars().all()), total @classmethod async def insert_activity(cls, session: AsyncSession, fields: dict): row = cls(**fields) session.add(row) await session.commit() return await cls.get_activity_by_id(session, row.id) @classmethod async def update_activity(cls, session: AsyncSession, record_id, fields: dict): row = await cls.get_activity_by_id(session, record_id) if not row: return None for key, value in fields.items(): setattr(row, key, value) session.add(row) await session.commit() await session.refresh(row) return row class Feedback(SQLModel, table=True): __tablename__ = "feedback" id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True) review: str = Field(default="") financial_status: str = Field(default="") score: float = Field(default=0.0) note: str | None = Field(default=None) created_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) updated_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) reviewed_by: uuid.UUID | None = Field(default=None, foreign_key="users.id") inbox_id: int | None = Field(default=None, foreign_key="inbox.id") user: Optional["Users"] = Relationship( back_populates="feedback", sa_relationship_kwargs={"lazy": "selectin"}, ) inbox: Optional["Inbox"] = Relationship( back_populates="feedback", sa_relationship_kwargs={"lazy": "selectin"}, ) @staticmethod def _as_uuid(record_id) -> uuid.UUID | None: try: return uuid.UUID(str(record_id)) except ValueError: return None @classmethod async def get_feedback_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_feedback_by_inbox(cls, session: AsyncSession, inbox_id: int): result = await session.execute( select(cls).where(cls.inbox_id == inbox_id).order_by(cls.created_at.desc()) ) return result.scalars().all() @classmethod async def insert_feedback(cls, session: AsyncSession, fields: dict): row = cls(**fields) session.add(row) await session.commit() return await cls.get_feedback_by_id(session, row.id) @classmethod async def update_feedback(cls, session: AsyncSession, record_id, fields: dict): row = await cls.get_feedback_by_id(session, record_id) if not row: return None for key, value in fields.items(): setattr(row, key, value) row.updated_at = _now() session.add(row) await session.commit() await session.refresh(row) return row class ApplicationStageTransitions(SQLModel, table=True): """Temporal history of application stage changes. Inbox moves write inbox_messages.application_status; manual-upload moves write manual_upload_candidate.status. Exactly one of inbox_id / manual_upload_candidate_id is set. NULL valid_to means the stage is current. """ __tablename__ = "application_stage_transitions" id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True) inbox_id: int | None = Field(default=None, index=True, foreign_key="inbox.id") manual_upload_candidate_id: uuid.UUID | None = Field( default=None, index=True, foreign_key="manual_upload_candidate.id" ) from_stage: str | None = Field(default=None) to_stage: str valid_from: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True)) valid_to: datetime | None = Field(default=None, sa_type=DateTime(timezone=True)) changed_by: uuid.UUID | None = Field(default=None, foreign_key="users.id") actor_kind: str = Field(default="user") change_reason: str | None = Field(default=None) created_at: datetime = Field(default_factory=_now, 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 fetch_by_inbox(cls, session: AsyncSession, inbox_id: int): result = await session.execute( select(cls).where(cls.inbox_id == int(inbox_id)).order_by(cls.valid_from.desc()) ) return list(result.scalars().all()) @classmethod async def fetch_by_manual(cls, session: AsyncSession, record_id): uid = cls._as_uuid(record_id) if uid is None: return [] result = await session.execute( select(cls).where(cls.manual_upload_candidate_id == uid).order_by(cls.valid_from.desc()) ) return list(result.scalars().all()) @classmethod async def get_open_transition(cls, session: AsyncSession, inbox_id=None, manual_upload_candidate_id=None): statement = select(cls).where(cls.valid_to.is_(None)).order_by(cls.valid_from.desc()) if inbox_id is not None: statement = statement.where(cls.inbox_id == int(inbox_id)) elif manual_upload_candidate_id is not None: uid = cls._as_uuid(manual_upload_candidate_id) if uid is None: return None statement = statement.where(cls.manual_upload_candidate_id == uid) else: return None result = await session.execute(statement) return result.scalars().first() @classmethod async def insert_transition(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 close_open(cls, session: AsyncSession, inbox_id=None, *, manual_upload_candidate_id=None, at: datetime | None = None, commit: bool = False): row = await cls.get_open_transition( session, inbox_id=inbox_id, manual_upload_candidate_id=manual_upload_candidate_id ) if not row: return None row.valid_to = at or _now() session.add(row) if commit: await session.commit() await session.refresh(row) return row @classmethod async def count_by_inbox(cls, session: AsyncSession, inbox_id: int): statement = select(func.count()).select_from(cls).where(cls.inbox_id == int(inbox_id)) result = await session.execute(statement) return result.scalar_one() import users.models as _users_models # noqa: E402, F401