import logging from typing import List, Optional import pandas as pd from fastapi import APIRouter, Body from service.api import schemas from service.tasks.api_service_handler.klaviyo.service_klaviyo import KlaviyoServiceCore from service.tasks.api_service_handler.klaviyo.service_klaviyo_db_handler import ( DbKlaviyoServiceHandler, ) from service.tasks.api_service_handler.utils import field_mapper from service.utils.appsync_communication.message_types import ServiceUserTokenError logger = logging.getLogger(__name__) router = APIRouter(tags=["Klaviyo"]) @router.post("/get_klaviyo_status", response_model=schemas.KlaviyoStatus) def get_klaviyo_status( management_schema: str = Body(..., embed=True), ): db_io = DbKlaviyoServiceHandler() df = db_io.read_from_klaviyo_user_table(schema=management_schema) if len(df) > 0: return {"status": "connected"} else: return {"status": "not_connected"} @router.post("/update_klaviyo_status", response_model=schemas.KlaviyoStatus) def update_klaviyo_status( user_id: str = Body(...), status: str = Body(...), management_schema: str = Body(...), publicKey: Optional[str] = Body(None), privateKey: Optional[str] = Body(None), ): db_io = DbKlaviyoServiceHandler() if status == "connected": if privateKey is None: raise RuntimeError("Please provide Private key") db_io.write_to_klaviyo_user_table( schema=management_schema, public_key=publicKey, private_key=privateKey, ) try: """Checking the keys before finalizing the process""" klaviyo = KlaviyoServiceCore(schema=management_schema) data = klaviyo.get_lists_and_segments() if data["data"].get("status", "") == 403: db_io.delete_from_klaviyo_user_table(schema=management_schema) raise ServiceUserTokenError( "Seems that provided keys are incorrect - please provide different keys" ) except Exception as e: logger.exception(f"Klaviyo user delete error: {e}") db_io.delete_from_klaviyo_user_table(schema=management_schema) raise ServiceUserTokenError( "Seems that provided keys are incorrect - please provide different keys" ) result = {"status": "connected"} elif status == "disconnected": db_io.delete_from_klaviyo_user_table(schema=management_schema) result = {"status": "not_connected"} else: result = {"status": "unknown"} return result @router.post("/list_klaviyo_lists", response_model=List[schemas.KlaviyoList]) @field_mapper def list_klaviyo_lists(management_schema: str = Body(..., embed=True)): db_io = DbKlaviyoServiceHandler() klaviyo = KlaviyoServiceCore(schema=management_schema) data = klaviyo.get_lists_and_segments() df = pd.DataFrame(data["data"]["data"]) df.rename( columns={"updated": "timestamp_updated", "created": "timestamp_created"}, inplace=True, ) df.drop(columns=["object"], inplace=True) df.sort_values(by="timestamp_updated", ascending=False, inplace=True) df.reset_index(inplace=True, drop=True) records = df.to_dict("records") db_io.delete_from_klaviyo_list_table(schema=management_schema) db_io.write_to_klaviyo_list_table(schema=management_schema, data=records) return records