""" Hive Task Model. This HiveTask model is used to store information about hive_task table. """ import uuid from typing import Any from sqlalchemy import ( Column, Integer, String, delete, ) from sqlalchemy.orm.session import Session from assets.connectors import mysql class HiveTask(mysql.AuModel): """Table definition for hive_task table.""" __tablename__ = "hive_task" id = Column(Integer, primary_key=True, autoincrement=True) task_id = Column(String, nullable=False) asset_final_id = Column(Integer, nullable=False) model = Column(String, nullable=False) model_version = Column(Integer, nullable=False) def as_dict(self) -> dict[str, Any]: """Return object as dict. Returns: dict: Dictionary representation of object """ return { "id": self.id, "task_id": self.task_id, "asset_final_id": self.asset_final_id, "model": self.model, "model_version": self.model_version, } def _create_hive_task( session: Session, asset_final_id: int, task: dict[str, Any] ) -> int | None: """Create hive_task item in the table. Args: session (Session): SQLAlchemy database session. asset_final_id (int): Asset final id. task (dict): task_id, model and model_version Returns: id (int): Last inserted record id. """ hive_task = HiveTask( asset_final_id=asset_final_id, task_id=str(uuid.UUID(task["task_id"])), model=task["model"], model_version=task["model_version"], ) session.add(hive_task) session.flush() return hive_task.id def _delete_hive_task(session: Session, asset_final_id: int) -> None: """Delete hive_task rows for an asset_final_id (used by the overwrite re-scan path). Args: session (Session): SQLAlchemy database session. asset_final_id (int): Asset final id. """ session.execute(delete(HiveTask).where(HiveTask.asset_final_id == asset_final_id))