from __future__ import annotations import json from fastapi import FastAPI from sqlalchemy import select from services.shared.core import new_id, utc_now_iso from services.shared.db import get_session from services.shared.models import EmailMessageIn, EmailMessageOut, HealthResponse from services.shared.sql_init import init_sql_schema from services.shared.sql_models import EmailMessageRow, Interaction, InteractionTimeline app = FastAPI(title="email-adapter-service", version="2.0.0") init_sql_schema() def _to_out(row: EmailMessageRow) -> EmailMessageOut: return EmailMessageOut( message_id=row.message_id, from_email=row.from_email, subject=row.subject, body=row.body, customer_external_id=row.customer_external_id, queue_id=row.queue_id, priority=row.priority, payload=json.loads(row.payload_json or "{}"), interaction_id=row.interaction_id, created_at=row.created_at, ) def _create_interaction(session, payload: EmailMessageIn, created_at: str) -> str: interaction_id = new_id("int") session.add( Interaction( interaction_id=interaction_id, channel="email", subject=payload.subject[:120], customer_id=payload.customer_external_id, queue_id=payload.queue_id or "q_email", priority=payload.priority, status="new", assigned_to=None, created_at=created_at, updated_at=created_at, ) ) session.add( InteractionTimeline( interaction_id=interaction_id, timestamp=created_at, action="interaction.created", metadata_json=json.dumps({"channel": "email"}, ensure_ascii=False), ) ) session.add( InteractionTimeline( interaction_id=interaction_id, timestamp=created_at, action="email.message_received", metadata_json=json.dumps({"from_email": payload.from_email}, ensure_ascii=False), ) ) return interaction_id @app.get("/health", response_model=HealthResponse) def health() -> HealthResponse: return HealthResponse(status="ok", service="email-adapter-service", version="v2") @app.post("/integrations/email/messages", response_model=EmailMessageOut) def create_message(payload: EmailMessageIn) -> EmailMessageOut: session = get_session() try: created_at = utc_now_iso() interaction_id = _create_interaction(session, payload, created_at) row = EmailMessageRow( message_id=new_id("eml"), from_email=payload.from_email, subject=payload.subject, body=payload.body, customer_external_id=payload.customer_external_id, interaction_id=interaction_id, queue_id=payload.queue_id or "q_email", priority=payload.priority, payload_json=json.dumps(payload.payload, ensure_ascii=False), created_at=created_at, ) session.add(row) session.commit() session.refresh(row) return _to_out(row) finally: session.close() @app.get("/integrations/email/messages", response_model=list[EmailMessageOut]) def list_messages(from_email: str | None = None, limit: int = 100) -> list[EmailMessageOut]: session = get_session() try: stmt = select(EmailMessageRow).order_by(EmailMessageRow.id.asc()) if from_email: stmt = stmt.where(EmailMessageRow.from_email == from_email) rows = session.execute(stmt).scalars().all() rows = rows[-limit:] return [_to_out(row) for row in rows] finally: session.close()