"""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)