"""Application DTO layer.""" from enum import Enum from typing import Dict from typing import List from typing import Optional from kafka.admin import NewTopic from pydantic import BaseModel from pydantic import Field from dbdeploy.util.common import str_uuid from dbdeploy.util.exceptions import KafkaDLQError # Environment information class EnvType(str, Enum): """ Enum class for environment types """ PROD = "prod" QA = "qa" DEV = "dev" TEST = "test" class Column(BaseModel): """DTO to describe mysql column.""" field: str type: str optional: bool = True class TableSchema(BaseModel): """DTO to describe mysql table schema.""" table_name: str primary_key: str fields: List[Column] optional: bool = True type: str = 'struct' class JsonSchema(BaseModel): """DTO to describe mysql sink record schema.""" fields: List[Column] optional: bool = False type: str = 'struct' class ChangeSet(BaseModel): """DTO to store changeset parameters.""" changeset_id: str precondition: bool = False sql_query: Optional[str] kafka_message: Optional[object | str] kafka_message_key: Optional[object | str] kafka_message_key_serializer: Optional[str] kafka_message_value_serializer: Optional[str] cypher_query: Optional[str] table_schema: Optional[TableSchema] = None run_on_change: bool = False run_always: bool = False skip_sink_connector: bool = False topic_name: Optional[str] snowflake_account: str neo4j_server: Optional[str] kafka_cluster_name: Optional[str] = 'managed-kafka-cdc-destination' insert_mode: Optional[str] = 'update' class ExecutionJob(ChangeSet): """DTO to store job parameters.""" id: str = Field(default_factory=str_uuid) topic: NewTopic = Field(default=NewTopic( name=None, num_partitions=1, replication_factor=3)) dlq_topic: NewTopic = Field(default=NewTopic( name=None, num_partitions=1, replication_factor=3)) connector_name: str = Field(default=None) offsets: Dict[int, int] = Field(default={}) is_done: bool = False topics_created: bool = False connector_started: bool = False main_error: Exception = Field(default=None) dlq_error: KafkaDLQError = Field(default=None) class Config: arbitrary_types_allowed = True class KafkaCluster(BaseModel): """Kafka Cluster DTO which holds cluster params.""" name: str bootstrap_brokers: str security_protocol: str = Field(default='SSL')