import logging import os import uuid from datetime import datetime, timezone from typing import Any, Optional from dotenv import load_dotenv from fastapi import HTTPException from inbox.enums import Candidate_application_Status from role.models import EnumRoles, Roles from sqlalchemy import Column, DateTime, func, or_, update from sqlalchemy.dialects.postgresql import JSONB from sqlalchemy.exc import IntegrityError from sqlalchemy.ext.asyncio import AsyncSession from sqlmodel import Field, Relationship, SQLModel, select, true from users.models import Users from users.plugins import hash_password load_dotenv() logger = logging.getLogger("inbox.models") # Placeholder only. The account lands inactive and the candidate is mailed a # confirmation link; the real password comes from the reset flow afterwards. DEFAULT_CANDIDATE_PASSWORD = os.getenv("DEFAULT_CANDIDATE_PASSWORD", "Utopia!@#") CANDIDATE_ROLE_ID_FALLBACK = 8 # mirrors users/views.py:signup_user SKIP_SENDER_PREFIXES = ("noreply", "no-reply", "donotreply", "do-not-reply", "mailer-daemon", "postmaster", "bounce") class Inbox(SQLModel, table=True): __tablename__ = "inbox" id: int | None = Field(default=None, primary_key=True) user_id: uuid.UUID | None = Field(default=None, foreign_key="users.id") alert_id: uuid.UUID | None = Field(default=None, foreign_key="inbox_alerts.id") alerts: Optional["Inbox_Alerts"] = Relationship(back_populates="inbox") message_id: uuid.UUID | None = Field(default=None, foreign_key="inbox_messages.id") messages: Optional["Inbox_Messages"] = Relationship(back_populates="inbox") created_at: datetime = Field(default_factory=datetime.now) updated_at: datetime = Field(default_factory=datetime.now) # is_active: bool = Field(default=True) # is_deleted: bool = Field(default=False) # user: Users | None = Relationship(back_populates="inbox") class Inbox_Alerts(SQLModel, table=True): __tablename__ = "inbox_alerts" id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True) alert_sender_name: str alert_sender_email: str is_read: bool = Field(default=False) recieve_time: datetime = Field(default_factory=datetime.now) inbox: list[Inbox] = Relationship(back_populates="alerts") class Inbox_Messages(SQLModel, table=True): __tablename__ = "inbox_messages" id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True) message_id: str | None = Field(default=None, index=True, unique=True) full_email_response: dict[str, Any] | None = Field( default=None, sa_column=Column(JSONB) ) application_status: Candidate_application_Status = Field(default=Candidate_application_Status.CLOSED) message_subject: str message_body: str message_sent_time: str message_received_time: str message_from: str message_to: str message_cc: str | None = Field(default=None) message_bcc: str | None = Field(default=None) message_read: bool = Field(default=False) attachment: bool = Field(default=False) message_reply: str | None = Field(default=None) file_name: str | None = Field(default=None) file_path: str | None = Field(default=None) resume_text: str | None = Field(default=None) experience: str | None = Field(default=None) suggested_job_post_ids: list[str] | None = Field(default=None, sa_column=Column(JSONB)) match_summary: str | None = Field(default=None) match_reasoning: str | None = Field(default=None) match_status: str | None = Field(default=None) match_error: str | None = Field(default=None) matched_at: datetime | None = Field(default=None, sa_type=DateTime(timezone=True)) inbox: list[Inbox] = Relationship(back_populates="messages") @staticmethod def _body_text(email_data: dict) -> str: body = email_data.get("body") if isinstance(body, dict): return body.get("content") or "" if isinstance(body, str): return body return email_data.get("bodyPreview") or "" # @classmethod # async def get_candidate_profile(cls,session:AsyncSession,user_id:uuid.UUID|None=None): # try: # qryy=select(cls,Users).join(cls,cls.) # if user_id # except Exception as e: # raise HTTPException(status_code=500,detail=str(e)) @classmethod async def get_all_applications(cls,session:AsyncSession,message_id:uuid.UUID|int|None=None): try: qry=select(cls.message_id,cls.full_email_response,cls.message_subject,cls.message_from,cls.message_to,cls.message_sent_time,cls.message_read,cls.attachment) if message_id: qry=qry.where(cls.message_id==message_id) result=await session.execute(qry) return result.scalars().all() except Exception as e: raise HTTPException(status_code=500,detail=str(e)) @classmethod async def set_match_result( cls, session: AsyncSession, record_id, *, resume_text=None, experience=None, suggested_job_post_ids=None, summary="", reasoning="", status="", error="", ): """Persist agent output onto one inbox row; returns the row or None.""" row = await cls.get_inbox_message_by_id(session, record_id) if not row: return None if resume_text is not None: row.resume_text = resume_text row.suggested_job_post_ids = suggested_job_post_ids row.match_summary = summary or None row.match_reasoning = reasoning or None row.match_status = status or None row.match_error = error or None row.experience = experience or None row.matched_at = datetime.now(timezone.utc) session.add(row) await session.commit() await session.refresh(row) return row @classmethod def _fields_from_email(cls, email_data: dict, file_path: list[str] | None = None) -> dict: return { "message_subject": email_data.get("subject") or "", "message_body": cls._body_text(email_data), "message_sent_time": email_data.get("sentDateTime") or "", "message_read": bool(email_data.get("isRead")), "message_received_time": email_data.get("receivedDateTime") or "", "message_from": email_data.get("from", {}) .get("emailAddress", {}) .get("address", ""), "message_to": ",".join( [r["emailAddress"]["address"] for r in email_data.get("toRecipients", [])] ), "message_cc": ",".join( [r["emailAddress"]["address"] for r in email_data.get("ccRecipients", [])] ), "message_bcc": ",".join( [r["emailAddress"]["address"] for r in email_data.get("bccRecipients", [])] ), "attachment": bool(email_data.get("hasAttachments")), "message_reply": ",".join( [r["emailAddress"]["address"] for r in email_data.get("replyTo", [])] ), "message_id": email_data.get("id"), "file_name": ",".join( [r.get("name") for r in email_data.get("attachments", [])] ), "file_path": ",".join(file_path) if file_path else None, "full_email_response": email_data, } @classmethod def _sender_address(cls, email_data: dict) -> str: return ( email_data.get("from", {}) .get("emailAddress", {}) .get("address", "") or "" ).strip().lower() @classmethod def _sender_display_name(cls, email_data: dict, address: str) -> str: name = ( email_data.get("from", {}) .get("emailAddress", {}) .get("name") or "" ).strip() if name: return name return address.split("@", 1)[0] if address else "candidate" @classmethod def _is_linkable_sender(cls, address: str) -> bool: if not address or "@" not in address: return False local = address.split("@", 1)[0] return not local.startswith(SKIP_SENDER_PREFIXES) @classmethod async def _link_sender(cls,session:AsyncSession,email_data:dict,email): address=cls._sender_address(email_data) if not cls._is_linkable_sender(address): return None try: user=(await session.execute( select(Users).where(func.lower(Users.email)==address) )).scalars().first() if not user: role=await Roles.get_role_by_name(session,EnumRoles.CANDIDATE.value) user=Users( name=cls._sender_display_name(email_data,address), email=address, role_id=role.id if role else CANDIDATE_ROLE_ID_FALLBACK, password=hash_password(DEFAULT_CANDIDATE_PASSWORD), ) session.add(user) session.add(Inbox(user_id=user.id,message_id=email.id)) await session.commit() return address link=(await session.execute( select(Inbox).where(Inbox.message_id==email.id,Inbox.user_id==user.id) )).scalars().first() if not link: session.add(Inbox(user_id=user.id,message_id=email.id)) await session.commit() return None except IntegrityError: await session.rollback() return None except Exception as e: await session.rollback() logger.warning("sender link failed for %s: %s",address,e) return None @classmethod async def insert_email( cls, session: AsyncSession, email_data: dict, file_path: list[str] | None = None, ): """Returns (row, new_user_email). new_user_email is set only when this call created the sender's Users row.""" fields = cls._fields_from_email(email_data, file_path) external_id = fields.get("message_id") link_user=None if external_id: existing = ( await session.execute( select(cls).where(cls.message_id == external_id) ) ).scalars().first() if existing: for key, value in fields.items(): setattr(existing, key, value) session.add(existing) await session.commit() await session.refresh(existing) if fields.get("attachment"): link_user=await cls._link_sender(session, email_data, existing) return existing, link_user email = cls(**fields) session.add(email) await session.commit() if fields.get("attachment"): link_user=await cls._link_sender(session, email_data, email) return email, link_user @classmethod def _search_filter(cls, search: str): pattern = f"%{search}%" return or_( cls.message_subject.ilike(pattern), cls.message_from.ilike(pattern), cls.message_body.ilike(pattern), ) @classmethod async def get_inbox_messages( cls, session: AsyncSession, top: int | None, skip: int, search: str | None, isread: bool=True, application_status: Candidate_application_Status=Candidate_application_Status.CLOSED ): statement = select(cls).order_by(cls.message_received_time.desc()) if search: statement = statement.where(cls._search_filter(search)) if application_status == Candidate_application_Status.PROCESS or application_status==Candidate_application_Status.REJECTED: statement = statement.where(cls.application_status==application_status) if skip: statement = statement.offset(skip) if top is not None: statement = statement.limit(top) if isread==False: statement = statement.where(cls.message_read==False) result = await session.execute(statement) return result.scalars().all() @classmethod async def get_inbox_message_by_id(cls, session: AsyncSession, record_id: str): try: uid = uuid.UUID(str(record_id)) except ValueError: return None result = await session.execute(select(cls).where(cls.id == uid)) return result.scalars().first() @classmethod async def count_inbox_messages(cls, session: AsyncSession, search: str | None, isread: bool=True, application_status: Candidate_application_Status=Candidate_application_Status.CLOSED): statement = select(func.count()).select_from(cls) if search: statement = statement.where(cls._search_filter(search)) if application_status == Candidate_application_Status.PROCESS or application_status==Candidate_application_Status.REJECTED: statement = statement.where(cls.application_status==application_status) if isread==False: statement = statement.where(cls.message_read==False) result = await session.execute(statement) return result.scalar_one() @classmethod async def apply_read_status(cls, session: AsyncSession, changes) -> int: """[{id, isRead, ...}] -> bulk UPDATE message_read. Returns rows touched. read is a ONE-WAY LATCH: only false -> true is applied, never the reverse. mark_message_read writes the local column only — nothing pushes the state back to Outlook — so upstream keeps reporting isRead=false and the every-minute sync_read_status sweep would otherwise revert a mail the user just opened. Cost of the latch: un-reading a mail in Outlook no longer propagates here. """ if not changes: return 0 read_ids=[c.get("id") for c in changes if c.get("id") and c.get("isRead")] if not read_ids: return 0 result=await session.execute( update(cls).where(cls.message_id.in_(read_ids)).values(message_read=True) ) await session.commit() return result.rowcount or 0 @classmethod async def mark_message_read(cls, session: AsyncSession, record_id): row=await cls.get_inbox_message_by_id(session,record_id) if not row: return None row.message_read=True session.add(row) await session.commit() await session.refresh(row) return row