"""Audio/text embedding application on modal."""

import os
import time
import modal
from suno_utils.worker.schema import QueueItem
import torch
import json
import requests
import shutil

from suno_utils.audio import Audio
from suno_utils.worker.loader import S3Loader
from suno_utils.worker.modal_base import MODAL_MOUNTS
from suno_utils.worker.utils import print_gpu_memory_usage
from suno_utils.worker.settings import s3_client
from suno_utils.gpt import chirp_v2_5 as chirp_v3
from suno_utils.models.ditto.ditto import Ditto
from suno_utils.worker.modal_base import get_modal_base_image

############## CHANGE THESE ##############

DEPLOYMENT_TYPE = "dev"  # dev, prod

##########################################


ENCODER_CONCURRENCY_LIMITS = {
    "dev": 200,
    "prod": 250,
}
KEEP_WARM = {
    "dev": 1,
    "prod": 1,
}

DITTO_EMBEDDING_DIM = 128

# set number of cpus.
VERBOSE_MESSAGE = DEPLOYMENT_TYPE == "dev"
MOUNT_PATH = "/suno/models"


aws_secret = modal.Secret.from_name("studio-aws")
SECRETS = [
    aws_secret,
    modal.Secret.from_dict(
        {
            "SUNO_ASSETS_PATH": "/suno/models/assets",
            "XDG_CACHE_HOME": "/suno/models/",
        }
    ),
    modal.Secret.from_name("openai-secret"),
    modal.Secret.from_name("turbopuffer-api-key"),
    modal.Secret.from_name("api-callback-token"),
]

# fix transformers for mert
base_image = get_modal_base_image().pip_install("transformers==4.44.0", "turbopuffer")

DITTO_S3_PATH = "s3://suno-data/victor/checkpoints/ditto/step_370k.pt"


class DittoWorker(S3Loader):
    def __init__(self):
        S3Loader.__init__(self)
        print("Start loading models")
        self.ditto_path = os.path.join(MOUNT_PATH, "ditto_models", "ditto.pt")
        self.music_encoder_path = os.path.join(MOUNT_PATH, "ditto_models", "musicfm.pt")
        if not os.path.exists(self.ditto_path):
            self.ditto_path = chirp_v3._get_model_if_needed(DITTO_S3_PATH, cache_dir=MOUNT_PATH)
        if not os.path.exists(self.music_encoder_path):
            self.music_encoder_path = chirp_v3._get_model_if_needed(
                "s3://suno-data/victor/checkpoints/ditto/musicfm_concat_epoch=51.pt",
                cache_dir=MOUNT_PATH,
            )
        print("found models locally.")
        self.ditto = Ditto(
            music_encoder_name="musicfm_concat",
            latent_dim=DITTO_EMBEDDING_DIM,
            model_path=self.ditto_path,
            music_encoder_path=self.music_encoder_path,
            is_flash=False,
        )

        self.ditto = self.ditto.eval().cuda()
        print("Finish loading models")

    @staticmethod
    def download_models(dir_path=MOUNT_PATH):
        print("Start downloading models")
        dl_path = chirp_v3._get_model_if_needed(DITTO_S3_PATH, cache_dir=dir_path)
        music_encoder_path = chirp_v3._get_model_if_needed(
            "s3://suno-data/victor/checkpoints/ditto/musicfm_concat_epoch=51.pt",
            cache_dir=dir_path,
        )
        print(f"Downloaded model to {dl_path}")
        print(f"Downloaded model to {music_encoder_path}")

        # Define source and destination paths
        source_ditto_path = dl_path
        source_music_encoder_path = music_encoder_path
        dest_dir = os.path.join(dir_path, "ditto_models")

        # Create destination directory if it doesn't exist
        os.makedirs(dest_dir, exist_ok=True)

        # Move Ditto model file
        dest_ditto_path = os.path.join(dest_dir, "ditto.pt")
        shutil.move(source_ditto_path, dest_ditto_path)
        print(f"Moved Ditto model to {dest_ditto_path}")

        # Move Music Encoder model file
        dest_music_encoder_path = os.path.join(dest_dir, "musicfm.pt")
        shutil.move(source_music_encoder_path, dest_music_encoder_path)
        print(f"Moved Music Encoder model to {dest_music_encoder_path}")
        print("Finish downloading models")


def download_model_wrapper_g():
    # this print is necessary to have modal rerun this when MODEL changes
    # Modal tracks referenced global variables
    # Change the name of the function to force a rerun
    print("Downloading model", DITTO_S3_PATH)
    DittoWorker.download_models()


image = base_image.run_function(download_model_wrapper_g, secrets=SECRETS)
APP_NAME = f"ditto-{DEPLOYMENT_TYPE}"
app = modal.App(APP_NAME, image=image)


@app.cls(
    gpu="T4",
    secrets=SECRETS,
    timeout=100,
    scaledown_window=400,
    mounts=MODAL_MOUNTS,
    retries=modal.Retries(
        max_retries=1,
        backoff_coefficient=2.0,
        initial_delay=5.0,
    ),
    max_containers=ENCODER_CONCURRENCY_LIMITS[DEPLOYMENT_TYPE],
    allow_concurrent_inputs=32,  # max 10 has ~ 5 GB at peak
    min_containers=0,  # TODO: split this into multiple apps
)
class DittoWorkerStub:
    def __init__(self):
        import torch

        num_gpus = torch.cuda.device_count()
        print(f"Found {num_gpus} GPUs.")

        self.worker = DittoWorker()

    @modal.method()
    def encode_audio(
        self,
        id: str = None,
        s3_id: str = None,
        s3_url: str = None,
        start: float = 0,
        dur: float = 30,
        callback_url=None,
    ) -> None:
        if VERBOSE_MESSAGE:
            print_gpu_memory_usage(self.__class__.__name__)
            torch.cuda.reset_max_memory_allocated()

        if dur > 30:
            raise ValueError("Duration must be less than 30 seconds")

        assert id or s3_id or s3_url, "Either id or s3_id or s3_url must be provided"

        if s3_id is None:
            s3_id = id
        elif s3_id != id:
            print(f"Warning: id and s3_id are different. Using s3_id {s3_id} for download.")

        if s3_id is not None:
            s3_url = f"s3://suno-data-uploads/studio/uploads/{s3_id}.mp3"

        try:
            s3_client.get_object(Bucket="suno-data-uploads", Key=f"studio/uploads/{s3_id}.mp3")
        except:
            print(f"File {id} doesn't exist on s3. Is it deleted?")
            return None

        audio = Audio.from_s3(s3_url, n_channels=1, sample_rate=24000)
        if audio.duration_s < 30:
            print(f"File {id} is too short for ditto encoding.")
            return None

        print(f"Starting ditto job with {id}, duration {audio.duration_s}.")
        audio = audio.get_segment(from_s=start, to_s=start + dur)
        wav = torch.tensor(audio.array_float).unsqueeze(0).cuda()
        emb = self.worker.ditto.music_to_latent(wav)[0].detach().cpu().numpy()

        if id and callback_url:
            requests.post(
                callback_url,
                json={"id": id, "vector": emb.tolist(), "index_name": "song-vector"},
                headers={
                    "Authentication": "Bearer 562a512f-0dce-4acd-bf23-ad9bb8d8a084",
                },
            )
        print(f"finished ditto job with {id}")
        return emb

    @modal.method()
    def encode_audio_and_upload(
        self, id: str = None, s3_url: str = None, start: float = 0, dur: float = 30
    ) -> None:
        """Use this method to backfill existing clips with embeddings.
        Should be used for one-off tasks (e.g. from your notebook) only.
        Please note that the DEPLOYMENT_TYPE should be set properly.
        For example, you can't upload staging clips to prod namespace"""
        if VERBOSE_MESSAGE:
            print_gpu_memory_usage(self.__class__.__name__)
            torch.cuda.reset_max_memory_allocated()

        if dur > 30:
            raise ValueError("Duration must be less than 30 seconds")

        assert id or s3_url, "Either id or s3_url must be provided"

        if id is not None:
            s3_url = f"s3://suno-data-uploads/studio/uploads/{id}.mp3"

        audio = Audio.from_s3(s3_url, n_channels=1, sample_rate=24000)
        audio = audio.get_segment(from_s=start, to_s=start + dur)
        wav = torch.tensor(audio.array_float).unsqueeze(0).cuda()
        emb = self.worker.ditto.music_to_latent(wav)[0].detach().cpu().numpy()

        import turbopuffer as tpuf

        tpuf.api_key = os.environ["TURBOPUFFER_API_KEY"]
        ns = tpuf.Namespace(f"energy-{DEPLOYMENT_TYPE}")
        ns.upsert(
            ids=[id],
            vectors=[emb.tolist()],
            distance_metric="cosine_distance",
        )

        return emb

    @modal.method()
    def encode_text(self, text: str) -> None:
        if VERBOSE_MESSAGE:
            print_gpu_memory_usage(self.__class__.__name__)
            torch.cuda.reset_max_memory_allocated()

        te = self.worker.ditto.text_to_latent("[CLS]" + text)[0].detach().cpu().numpy()
        return te

    @modal.method()
    def encode_texts(self, queue_item_json: str) -> list[dict]:
        queueItem = QueueItem(**json.loads(queue_item_json))
        metadata = queueItem.metadata
        texts = metadata.get("encode_texts_in_ditto", [])
        embeddings = []
        for text in texts:
            te = self.worker.ditto.text_to_latent("[CLS]" + text)[0].detach().cpu().numpy()
            embeddings.append({"text": text, "embedding": te.tolist()})
        queueItem.notify_progress({"id": queueItem.id, "embeddings": embeddings})
        return embeddings


@app.local_entrypoint()
def main():
    ditto_worker = DittoWorkerStub()
    print(ditto_worker.encode_audio.remote(id="4a77dea7-19f3-46d2-8b0a-b2b7e9ea9a05"))
    print(
        ditto_worker.encode_audio.remote(
            s3_url="s3://suno-data-uploads/studio/uploads/4a77dea7-19f3-46d2-8b0a-b2b7e9ea9a05.mp3"
        )
    )
    print(ditto_worker.encode_text.remote("jazz"))

    for i in range(2):
        ditto_worker.encode_audio.spawn("4a77dea7-19f3-46d2-8b0a-b2b7e9ea9a05")

    if DEPLOYMENT_TYPE == "dev":
        staging_clip_id = "094228d9-fbeb-4446-ae52-bb0c6343e35e"
        emb = ditto_worker.encode_audio_and_upload.spawn(id=staging_clip_id)
        assert len(emb.get()) == DITTO_EMBEDDING_DIM

    time.sleep(100)
