from typing import List from service.async_task_manager import io_task from service.tasks.api_service_handler.utils import ServiceRequestsHandler from .service_klaviyo_db_handler import DbKlaviyoServiceHandler class KlaviyoServiceCore(ServiceRequestsHandler): base_path = "https://a.klaviyo.com/api/" def __init__(self, schema): self.schema = schema self.public_key = None self.private_key = None self._get_klaviyo_credentials() def make_request(self, **kwargs): """Using existing params dictionary or creating a new one and inserting api_key value""" kwargs["params"] = kwargs.get("params", {}) kwargs["params"]["api_key"] = self.private_key """Converting version to correct path, once done removing the argument if exists""" version = kwargs.get("version", "v2") kwargs.pop("version", None) base_path = f"{self.base_path}{version}/" return super().make_request(base_path=base_path, **kwargs) def get_lists_and_segments(self): """Using different base path as could not get V2 to work (seems to be Klaviyo's bug)""" response = self.make_request(method="GET", endpoint="lists", version="v1") return response def _get_klaviyo_credentials(self): db_io = DbKlaviyoServiceHandler() creds = db_io.read_from_klaviyo_user_table(schema=self.schema) creds = creds.to_dict("records") if len(creds) > 0: try: self.public_key = creds[0]["public_key"] self.private_key = creds[0]["private_key"] except Exception: raise RuntimeError("Something wrong with stored Klaviyo Credentials") else: raise RuntimeError("No credentials Stored.") def get_list_or_segment_persons( self, list_id ): # primary_klaviyo_collection (parent_collection) response_list = [] endpoint = f"group/{list_id}/members/all" response = self.make_request(method="GET", endpoint=endpoint, version="v2") response_list.append(response) if "data" in response: data = response["data"] marker = data.get("marker", None) """marker is used to indicate next batch of users (1000) per call""" while marker is not None: response = self.make_request( method="GET", endpoint=endpoint, version="v2", params={"marker": marker}, ) response_list.append(response) marker = response["data"].get("marker", None) return response_list def get_person_static_data( self, person_id ): # enrich_klaviyo_static_data (separate_child_collection) endpoint = f"person/{person_id}" response = self.make_request(method="GET", endpoint=endpoint, version="v1") return response @io_task def get_person_static_data_async(self, person_id, task_handle): endpoint = f"person/{person_id}" response = self.make_request(method="GET", endpoint=endpoint, version="v1") return {"response": response, "person_id": person_id} def get_person_events_data( self, person_id, since=None ): # enrich_klaviyo_events_data (separate_child_collection) endpoint = f"person/{person_id}/metrics/timeline" params = {} if since is not None: params["since"] = since response = self.make_request( method="GET", endpoint=endpoint, version="v1", params=params ) return response @io_task def get_person_events_data_async(self, person_id, task_handle): """For async run we ensure that obtained all the data for person_id""" response_list = [] response = self.get_person_events_data(person_id=person_id) if "data" in response: data = response["data"] events = data.get("data", []) for event in events: event["klaviyo_person_id"] = person_id response_list.append(event) """Next variable would indicate that there is another page for the same user""" next_page = data.get("next", None) while next_page is not None: """Extracting all the event's data""" response = self.get_person_events_data( person_id=person_id, since=next_page ) data = response["data"] for event in data["data"]: event["klaviyo_person_id"] = person_id response_list.append(event) next_page = data.get("next", None) return response_list def create_list(self, list_name): """https://apidocs.klaviyo.com/reference/lists-segments#create-list""" data = f"list_name={list_name}" endpoint = "lists" headers = { "accept": "application/json", "content-type": "application/x-www-form-urlencoded", } response = self.make_request( method="POST", endpoint=endpoint, version="v2", data=data, headers=headers ) return response def _add_users_to_a_list(self, list_id, users: List[dict]): endpoint = f"list/{list_id}/members" data = {"profiles": users} response = self.make_request( method="POST", endpoint=endpoint, version="v2", json=data ) return response @io_task def add_users_to_a_list(self, list_id, users: List[dict], task_handle): return self._add_users_to_a_list(list_id, users)