""" TranscodingJob Model. This TranscodingJob model uses sqlalchemy. It's used to store information about transcoding orders. """ from datetime import datetime from typing import Any from owsresponse import response from sentry_sdk import capture_exception from sqlalchemy import Column, ColumnElement, Enum, String, Text, func from sqlalchemy.dialects.mysql import INTEGER, TIMESTAMP from sqlalchemy.exc import SQLAlchemyError from sqlalchemy.orm import Mapped from transcoding.connectors import mysql from transcoding.constants import error, transcoding class TranscodingJob(mysql.BaseModel): """Table definition for transcoding_job table. Transcoding job entry is responsible for storing information about status of transcoding (status, attempt, description, created_timestamp, updated_timestamp), transcdoing properties (container, channels, codec, sample_rate, bit_rate, bit_depth), and other values related to asset transcoding. Transcoding properties clarification: channels: Number of channels, e.g.: 1 or 2 codec: Codec which has been used for audio encoding, e.g.: pcm, .. container: Container format sample_rate: Number of samples per second in Hz, e.g.: 44100, 48000,.. bit_rate: Number of bits per second in bit/s, e.g.: 1411200, 1519552,. bit_depth: Number of bits of information in each sample e.g.: 16, 24 If these properties are not specified, we use the values from the source file. """ __tablename__ = "transcoding_job" transcoding_job_id = Column( INTEGER(unsigned=True), primary_key=True, autoincrement=True ) transcoding_order_id = Column(INTEGER(unsigned=True), nullable=False) status: Mapped[str] = Column( Enum(*transcoding.STATUSES), nullable=False, default=transcoding.REQUESTED_STATUS, ) description = Column(Text(), default=None, nullable=True) attempt = Column(INTEGER(display_width=2, unsigned=True), nullable=False, default=0) output_bucket = Column(String(255), default=None, nullable=True) output_key = Column(String(255), default=None, nullable=True) container = Column(String(255), default=None, nullable=True) channels = Column(INTEGER(display_width=1, unsigned=True), nullable=True) codec = Column(String(255), default=None, nullable=True) sample_rate = Column(INTEGER(11), default=None, nullable=True) bit_rate = Column(INTEGER(11), default=None, nullable=True) bit_depth = Column(INTEGER(11), default=None, nullable=True) duration = Column(INTEGER, default=None, nullable=True) created_timestamp: Mapped[datetime] = Column( TIMESTAMP, nullable=False, default=func.now() ) updated_timestamp: Mapped[datetime] = Column( TIMESTAMP, nullable=False, default=datetime.utcnow, onupdate=func.now() ) def as_dict(self) -> dict[str, Any]: """Return object as dict. Returns: dict: Dictionary representation of object. """ transcoding_job_dict = { "transcoding_job_id": self.transcoding_job_id, "transcoding_order_id": self.transcoding_order_id, "status": self.status, "description": self.description, "attempt": self.attempt, "output_bucket": self.output_bucket, "output_key": self.output_key, "container": self.container, "channels": self.channels, "codec": self.codec, "sample_rate": self.sample_rate, "bit_rate": self.bit_rate, "bit_depth": self.bit_depth, "duration": self.duration, "created_timestamp": self.created_timestamp.isoformat(), "updated_timestamp": self.updated_timestamp.isoformat(), } return transcoding_job_dict def create_transcoding_job( transcoding_order_id: int, fields: dict[str, Any] ) -> response.Response: """Create new transcoding job item in the table. Args: transcoding_order_id (int): Transcoding order id. fields (dict): optional values for transcoding job. Returns: response.Response: Created transcoding job info or error. """ insertion_values = {"transcoding_order_id": transcoding_order_id, **fields} try: transcoding_job = TranscodingJob(**insertion_values) with mysql.db_session() as session: session.add(transcoding_job) session.flush() transcoding_job_dict = transcoding_job.as_dict() return response.Response(transcoding_job_dict) except SQLAlchemyError as e: capture_exception(e) return response.create_fatal_response(e.args) def get_transcoding_job_by_id(transcoding_job_id: int) -> response.Response: """Return transcoding job info by id. Args: transcoding_job_id (int): Transcoding job id. Returns: response.Response: TranscodingJob.as_dict() in message attribute or error response. """ try: with mysql.db_session() as session: filters = [(TranscodingJob.transcoding_job_id == transcoding_job_id)] entry = session.query(TranscodingJob).filter(*filters).one_or_none() if entry is None: return response.create_not_found_response( message=error.ERROR_NOT_FOUND_MESSAGE.format("Transcoding Job") ) transcoding_job_dict = entry.as_dict() return response.Response(transcoding_job_dict) except SQLAlchemyError as e: capture_exception(e) return response.create_fatal_response(e.args) def get_transcoding_jobs_by_order_id(transcoding_order_id: int) -> response.Response: """Return transcoding jobs info by transcoding_order_id. Args: transcoding_order_id (int): Transcoding order id. Returns: response.Response: TranscodingJob.as_dict() in message attribute or error response. """ try: with mysql.db_session() as session: filters = [(TranscodingJob.transcoding_order_id == transcoding_order_id)] transcoding_jobs = session.query(TranscodingJob).filter(*filters).all() result_data = [entry.as_dict() for entry in transcoding_jobs] return response.Response(message=result_data) except SQLAlchemyError as e: capture_exception(e) return response.create_fatal_response(e.args) def update_transcoding_job_attempt_counter( transcoding_job_id: int, attempt: int ) -> response.Response: """Update attempt value in one transcoding job by provided filters. Args: transcoding_job_id (int): Existing transcoding job id. attempt (int): New attempt counter value. Returns: response.Response: TranscodingJob.as_dict() in message attribute or error response. """ filters = [(TranscodingJob.transcoding_job_id == transcoding_job_id)] updated_values = {"attempt": attempt} return _update_transcoding_job(filters, updated_values) def update_transcoding_job_status( transcoding_job_id: int, new_status: str, new_description: str | None = None, output_bucket: str | None = None, output_key: str | None = None, metadata: dict[str, Any] | None = None, ) -> response.Response: """Update status value in one transcoding job by provided filters. "metadata": { "@type": "Audio", "Format": "PCM", "Format_Settings_Endianness": "Little", "Format_Settings_Sign": "Signed", "CodecID": "1", "Duration": "2.963", "BitRate_Mode": "CBR", "BitRate": "1411200", "Channels": "2", "SamplingRate": "44100", "SamplingCount": "130688", "BitDepth": "16", "StreamSize": "522752" } Args: transcoding_job_id (int): Existing transcoding job id. new_status (int): New status value must be one of the following: ['requested', 'completed', 'processing', 'error'] new_description (str): Description of transcoding job processing. output_bucket (str): Transcoded file S3 bucket. output_key (str): Trancoded file name. metadata (dict): Transcoded file metadata. Returns: response.Response: TranscodingJob.as_dict() in message attribute or error response. """ filters = [(TranscodingJob.transcoding_job_id == transcoding_job_id)] updated_values: dict[str, Any] = {"status": new_status} if new_description is not None: updated_values["description"] = new_description if output_bucket is not None and output_key is not None: updated_values["output_bucket"] = output_bucket updated_values["output_key"] = output_key if metadata: updated_values["duration"] = int(float(metadata.get("Duration", 0)) * 1000) updated_values["channels"] = int(metadata.get("Channels", 0)) updated_values["codec"] = metadata.get("Format", "").lower() updated_values["sample_rate"] = int(metadata.get("SamplingRate", 0)) updated_values["bit_rate"] = int(metadata.get("BitRate", 0)) updated_values["bit_depth"] = int(metadata.get("BitDepth", 0)) return _update_transcoding_job(filters, updated_values) def _update_transcoding_job( filters: list[ColumnElement[Any]], updated_values: dict[Any, Any] ) -> response.Response: """Update one trancoding job row by provided filters. Args: filters (list): Values for filtering. updated_values (dict): New values for existed entry. Returns: response.Response: TranscodingJob.as_dict() in message attribute or error response. """ try: with mysql.db_session() as session: updated_items = ( session.query(TranscodingJob).filter(*filters).update(updated_values) ) if updated_items != 1: return response.create_fatal_response( error.ERROR_MESSAGE_TRANSCODING_JOB_NOT_UPDATED ) session.commit() row = session.query(TranscodingJob).filter(*filters).first() if row: return response.Response(row.as_dict()) else: return response.create_not_found_response() except SQLAlchemyError as e: capture_exception(e) return response.create_fatal_response(e.args)