167 lines
6.0 KiB
Python
167 lines
6.0 KiB
Python
import uuid
|
|
from datetime import datetime
|
|
from typing import Any, Optional
|
|
|
|
from sqlalchemy import Column, func, or_
|
|
from sqlalchemy.dialects.postgresql import JSONB
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
from sqlmodel import Field, Relationship, SQLModel, select
|
|
|
|
from users.models import Users
|
|
|
|
|
|
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)
|
|
)
|
|
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_path: str | None = Field(default=None)
|
|
|
|
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
|
|
def _fields_from_email(cls, email_data: dict) -> 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"),
|
|
"full_email_response": email_data,
|
|
}
|
|
|
|
@classmethod
|
|
async def insert_email(cls, session: AsyncSession, email_data: dict):
|
|
fields = cls._fields_from_email(email_data)
|
|
external_id = fields.get("message_id")
|
|
|
|
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)
|
|
return existing
|
|
|
|
email = cls(**fields)
|
|
session.add(email)
|
|
await session.commit()
|
|
return email
|
|
|
|
@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
|
|
):
|
|
statement = select(cls).order_by(cls.message_received_time.desc())
|
|
if search:
|
|
statement = statement.where(cls._search_filter(search))
|
|
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 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):
|
|
statement = select(func.count()).select_from(cls)
|
|
if search:
|
|
statement = statement.where(cls._search_filter(search))
|
|
result = await session.execute(statement)
|
|
return result.scalar_one()
|